diff --git a/Cargo.lock b/Cargo.lock index 025a00fdf..21b9512f4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -364,7 +364,7 @@ dependencies = [ [[package]] name = "asap-aware-mapping" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=344369e25c50e65dbaaff95b4ea628c31999cb59#344369e25c50e65dbaaff95b4ea628c31999cb59" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=a049a3b6c17cee4c51c79e478799b32280a7d71d#a049a3b6c17cee4c51c79e478799b32280a7d71d" dependencies = [ "asap-types", "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib)", @@ -376,7 +376,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=344369e25c50e65dbaaff95b4ea628c31999cb59#344369e25c50e65dbaaff95b4ea628c31999cb59" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=a049a3b6c17cee4c51c79e478799b32280a7d71d#a049a3b6c17cee4c51c79e478799b32280a7d71d" dependencies = [ "asap-types", "promql-parser 0.10.0 (git+https://github.com/ProjectASAP/promql-parser?rev=9fede7eecca923c9882fe256484d00d37f8706cb)", @@ -385,7 +385,7 @@ dependencies = [ [[package]] name = "asap-frontend-sql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=344369e25c50e65dbaaff95b4ea628c31999cb59#344369e25c50e65dbaaff95b4ea628c31999cb59" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=a049a3b6c17cee4c51c79e478799b32280a7d71d#a049a3b6c17cee4c51c79e478799b32280a7d71d" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -396,11 +396,12 @@ dependencies = [ [[package]] name = "asap-physical-operators" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=344369e25c50e65dbaaff95b4ea628c31999cb59#344369e25c50e65dbaaff95b4ea628c31999cb59" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=a049a3b6c17cee4c51c79e478799b32280a7d71d#a049a3b6c17cee4c51c79e478799b32280a7d71d" dependencies = [ "asap-types", "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib?rev=5f03ccbd798ed5fec62bdd839bcb331123cab369)", "futures", + "regex", "serde", "serde_json", "thiserror 2.0.20", @@ -410,12 +411,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=344369e25c50e65dbaaff95b4ea628c31999cb59#344369e25c50e65dbaaff95b4ea628c31999cb59" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=a049a3b6c17cee4c51c79e478799b32280a7d71d#a049a3b6c17cee4c51c79e478799b32280a7d71d" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=344369e25c50e65dbaaff95b4ea628c31999cb59#344369e25c50e65dbaaff95b4ea628c31999cb59" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=a049a3b6c17cee4c51c79e478799b32280a7d71d#a049a3b6c17cee4c51c79e478799b32280a7d71d" dependencies = [ "serde", "serde_json", diff --git a/Cargo.toml b/Cargo.toml index e31c43bf8..982dabb9e 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -16,10 +16,10 @@ version = "0.1.0" [workspace.dependencies] # Keep Planner frontends, selection, and IR on the same immutable revision. # Alias upstream asap-types because this workspace also defines asap_types. -planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "344369e25c50e65dbaaff95b4ea628c31999cb59" } -asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "344369e25c50e65dbaaff95b4ea628c31999cb59" } -asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "344369e25c50e65dbaaff95b4ea628c31999cb59" } -asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "344369e25c50e65dbaaff95b4ea628c31999cb59" } +planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "a049a3b6c17cee4c51c79e478799b32280a7d71d" } +asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "a049a3b6c17cee4c51c79e478799b32280a7d71d" } +asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "a049a3b6c17cee4c51c79e478799b32280a7d71d" } +asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "a049a3b6c17cee4c51c79e478799b32280a7d71d" } # Shared external deps (used by 2+ crates) serde = { version = "1.0", features = ["derive"] } @@ -39,7 +39,7 @@ arc-swap = "1.7" reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] } # Internal crates -asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "344369e25c50e65dbaaff95b4ea628c31999cb59" } +asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "a049a3b6c17cee4c51c79e478799b32280a7d71d" } asap_sketch_codec = { path = "crates/asap_sketch_codec" } asap_summary_state = { path = "crates/asap_summary_state" } asap_types = { path = "crates/asap_types" } diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 251856e99..14f1efd58 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -782,65 +782,8 @@ impl AccuracyEvidenceProvider for QueryEvidence<'_> { } } -/// Planner's row-binary error when per-series rows carry no series identity. -/// Planner reports it only as text, so reselection matches the message. -const SERIES_IDENTITY_REQUIRED: &str = - "row binary requires grouped rows or rows with a series identity"; - -/// Admit the workload selected over identity-typed roots as a candidate -/// forest, or give the reason it is rejected. -fn typed_reselection_forest( - typed: &mut [QueryCompilationInput], - canonical: &[QueryCompilationInput], - retyped: &[bool], - needs_identity: &dyn Fn(&Rc) -> bool, - inputs: &BackendLocalPhysicalInputs, - target: PhysicalDeploymentTarget, -) -> Result<(), String> { - if !canonical.iter().zip(&*typed).any(|(canonical, typed)| { - needs_identity(&canonical.selected_plan_root) && !needs_identity(&typed.selected_plan_root) - }) { - return Err("no computation compiles over identity-typed roots".into()); - } - for query in typed.iter_mut() { - prepare_window_implementations( - query, - &inputs.window_cost_model, - target, - inputs.query_retention_margin_ms, - ) - .map_err(|error| format!("{}: {error}", query.query_id))?; - } - // A root that cannot be retyped keeps its canonical states. Sharing one - // with a retyped root would give one deployed output two definitions. - let fingerprints = |retyped_side: bool| -> Result, String> { - let mut all = BTreeSet::new(); - for (query, _) in typed - .iter() - .zip(retyped) - .filter(|(_, &retyped)| retyped == retyped_side) - { - all.extend(state_fingerprints(query, target)?); - } - Ok(all) - }; - if !fingerprints(true)?.is_disjoint(&fingerprints(false)?) { - return Err("a root without a series identity shares state with a retyped root".into()); - } - // Native realizations of typed roots are bound only through their - // lifecycle placement. - for (query, &retyped) in typed.iter_mut().zip(retyped) { - if retyped { - query.retain(None) - } else { - query.retain_physical_candidate() - } - .map_err(|error| error.to_string())?; - } - Ok(()) -} - /// Deployed-output fingerprints of the states a query's selected root reads. +#[cfg(test)] fn state_fingerprints( query: &QueryCompilationInput, target: PhysicalDeploymentTarget, @@ -1082,90 +1025,51 @@ impl BackendLocalPlanningInput { exact_costs_by_id.insert(format!("compat-query-{index}"), rows.clone()); } } + // PromQL rows carry their complete identity before PlanSpace chooses + // shared states, so arithmetic never needs a second logical selection. + let logical_roots = canonical_roots + .iter() + .map(|root| { + planner_types::pre_asap::schema::with_promql_series_identity(root) + .map_or_else(|_| Rc::clone(root), Rc::new) + }) + .collect(); let mut planner_selection_trace = select_logical_roots_with_scoped_evidence_and_trace( &mut queries, - canonical_roots.clone(), + logical_roots, &topk_evidence_by_id, &scoped_evidence_by_id, &exact_costs_by_id, self.physical_inputs.erp.as_ref(), self.environment.observed_at_unix_ms, )?; - for query in &mut queries { + for (index, query) in queries.iter_mut().enumerate() { + // The preferred typed root can itself have a native realization; + // apply the same lifecycle pricing as the alternative forests. + if let Some(timed) = placement::time_native_candidate( + &query.selected_plan_root, + query, + &workload, + &data_workload, + index, + &self.environment, + &self.physical_inputs, + ) { + query.selected_plan_root = timed.root; + query.retain(Some(timed.physical))?; + planner_selection_trace.extend(timed.trace); + } prepare_window_implementations( query, &self.physical_inputs.window_cost_model, self.environment.target, self.physical_inputs.query_retention_margin_ms, )?; - query.retain_physical_candidate()?; - } - let mut planner_candidate_forests = Vec::new(); - // Planner matches per-series rows in query-time arithmetic only by the - // series identity. When a selected computation fails to compile for - // lack of it, the workload selected over identity-typed roots is one - // more forest; deployment pricing chooses among all forests. - let needs_identity = |root: &Rc| { - crate::query_plan::is_query_computation(root) - && crate::query_plan::compile_query_computation(root) - .is_err_and(|error| error.to_string().contains(SERIES_IDENTITY_REQUIRED)) - }; - if queries - .iter() - .any(|query| needs_identity(&query.selected_plan_root)) - { - let reject = |reason: String| { - serde_json::json!({ - "stage": "planner.series_identity_reselection", - "status": "rejected", "reason": reason, - }) - }; - let typed_roots: Vec<_> = canonical_roots - .iter() - .map(|root| { - asap_physical_operators::physical_planner::promql_rows::with_series_identity( - root, - ) - .map_or_else(|_| Rc::clone(root), Rc::new) - }) - .collect(); - let retyped: Vec<_> = typed_roots - .iter() - .zip(&canonical_roots) - .map(|(typed, canonical)| !Rc::ptr_eq(typed, canonical)) - .collect(); - let mut typed = queries.clone(); - match select_logical_roots_with_scoped_evidence_and_trace( - &mut typed, - typed_roots, - &topk_evidence_by_id, - &scoped_evidence_by_id, - &exact_costs_by_id, - self.physical_inputs.erp.as_ref(), - self.environment.observed_at_unix_ms, - ) { - Err(error) => planner_selection_trace.push(reject(error.to_string())), - Ok(trace) => { - planner_selection_trace.extend(trace.into_iter().map(|mut event| { - if let Some(event) = event.as_object_mut() { - event.insert("reselection".into(), "series_identity".into()); - } - event - })); - match typed_reselection_forest( - &mut typed, - &queries, - &retyped, - &needs_identity, - &self.physical_inputs, - self.environment.target, - ) { - Ok(()) => planner_candidate_forests.push(typed), - Err(reason) => planner_selection_trace.push(reject(reason)), - } - } + if query.physical_candidate.is_none() { + query.retain_physical_candidate()?; } } + let mut planner_candidate_forests = Vec::new(); // Native physical realizations need the complete series identity in // their rows. Planner's PlanSpace proposes them for the identity-typed // root; each is a logical alternative whose readout-built states are @@ -1219,6 +1123,7 @@ impl BackendLocalPlanningInput { &data_workload, index, &self.environment, + &self.physical_inputs, ) else { if typed_search { planner_selection_trace.push(serde_json::json!({ @@ -1416,7 +1321,8 @@ fn preserve_native_unsafe_raw_roots( /// Canonicalize Planner candidates that have no backend-maintained state at /// the physical compiler boundary. Their selected post-ASAP shape may be /// intentionally unsupported (and therefore invalid as an executable -/// maintenance DAG), but the original query remains a valid exact plan. +/// maintenance DAG), or a valid native program can lack any bindable stored +/// input (for example, an offset selector). Preserve the original exact query. fn preserve_invalid_exact_fallback_roots( queries: &mut [QueryCompilationInput], canonical_roots: &[Rc], @@ -1428,9 +1334,11 @@ fn preserve_invalid_exact_fallback_roots( query_id: query.query_id.clone(), reason, })?; - let invalid_executable = - selected.is_empty() && validate_executable_subdag(&query.selected_plan_root).is_err(); - if invalid_executable + let unbound_candidate = selected.is_empty() + && (validate_executable_subdag(&query.selected_plan_root).is_err() + || (query.physical_candidate.is_some() + && !super::maintained_population::supported_node(&query.selected_plan_root))); + if unbound_candidate && !matches!(query.selected_plan_root.expr, SummaryExpr::KeepPreAsap(_)) { let parsed = original_root(query, index, canonical_roots)?; @@ -1473,11 +1381,10 @@ fn preserve_uncompiled_computation_roots( Ok(()) } -/// A MetricsQL query whose only selected states are Prometheus-specific -/// counter readouts has no backend materialization to bind. Keep the original -/// query as one native exact root. Mixed queries retain their other selected -/// summaries and let query-time lowering cut only the counter branches. -fn preserve_metricsql_counter_only_roots( +/// MetricsQL counter readouts cannot bind Prometheus-specific stored states. +/// Preserve the whole native query when any branch requires such a readout, +/// including binary roots whose other operand could use a stored summary. +fn preserve_metricsql_counter_roots( queries: &mut [QueryCompilationInput], canonical_roots: &[Rc], composable: bool, @@ -1488,18 +1395,16 @@ fn preserve_metricsql_counter_only_roots( query_id: query.query_id.clone(), reason, })?; - if selected.is_empty() - || !selected.iter().all(|state| { - matches!( - state.family, - SummaryFamilyType::ExactAggregate( - planner_types::post_asap::ExactKind::Rate - | planner_types::post_asap::ExactKind::Increase, - _ - ) + if !selected.iter().any(|state| { + matches!( + state.family, + SummaryFamilyType::ExactAggregate( + planner_types::post_asap::ExactKind::Rate + | planner_types::post_asap::ExactKind::Increase, + _ ) - }) - { + ) + }) { continue; } let parsed = original_root(query, index, canonical_roots)?; @@ -1598,7 +1503,7 @@ impl DeploymentPlanCompiler { } if frontend == QueryFrontend::MetricsQl { - preserve_metricsql_counter_only_roots( + preserve_metricsql_counter_roots( &mut request.queries, &request.canonical_roots, request.allow_mixed_summary_and_exact_execution, @@ -1617,6 +1522,14 @@ impl DeploymentPlanCompiler { preserve_native_unsafe_raw_roots(&mut request.queries, &request.canonical_roots)?; } + // Native exact roots cannot retain a physical program whose stored inputs + // belonged to the candidate replaced by the frontend/source guards. + for query in &mut request.queries { + if matches!(query.selected_plan_root.expr, SummaryExpr::KeepPreAsap(_)) { + query.physical_candidate = None; + } + } + let roots = request .queries .iter() @@ -1627,11 +1540,6 @@ impl DeploymentPlanCompiler { request.queries[id].selected_plan_root = root; } let placement = placement::place(&request, &environment, frontend); - for (index, query) in request.queries.iter_mut().enumerate() { - if let Some(root) = placement.root(index) { - query.selected_plan_root = Rc::clone(root); - } - } let population_operators = super::maintained_population::operators(&request)?; let mut compiled_materializations = Vec::with_capacity(request.queries.len()); let mut output_computations = BTreeMap::<_, OutputComputation>::new(); @@ -5468,27 +5376,20 @@ pub(crate) mod tests { snapshot } - fn needs_series_identity(root: &Rc) -> bool { - crate::query_plan::compile_query_computation(root) - .is_err_and(|error| error.to_string().contains(SERIES_IDENTITY_REQUIRED)) - } - - // The identity-typed reselection is only a candidate forest: the primary - // workload stays the canonical selection, and its trace entries are tagged. + // Per-series arithmetic compiles on the first selection, without retrying PlanSpace. #[test] - fn typed_reselection_is_only_a_candidate_forest() { + fn logical_candidates_have_series_identity_before_selection() { let (request, _) = exact_workload_snapshot(&["avg_over_time(data[5m])"]) .into_physical_compilation_request() .unwrap(); - assert!(needs_series_identity( + assert!(crate::query_plan::compile_query_computation( &request.queries[0].selected_plan_root - )); - let typed = &request.planner_candidate_forests[0]; - assert!(crate::query_plan::compile_query_computation(&typed[0].selected_plan_root).is_ok()); - assert!(request - .planner_selection_trace - .iter() - .any(|event| event["reselection"] == "series_identity")); + ) + .is_ok()); + assert!(request.planner_selection_trace.iter().all(|event| { + event["reselection"].is_null() + && event["stage"] != "planner.series_identity_reselection" + })); } // Workloads mixing per-series arithmetic with aggregates of the same @@ -5526,8 +5427,7 @@ pub(crate) mod tests { } } - // A state shared by a retyped and a canonical root has the same deployed - // fingerprint either way, so it must keep one definition. + // Typing changes physical schemas without changing the stored-state identity. #[test] fn shared_state_fingerprint_is_independent_of_series_identity() { let (request, _) = @@ -5535,59 +5435,35 @@ pub(crate) mod tests { .into_physical_compilation_request() .unwrap(); let target = PhysicalDeploymentTarget::BackendLocalRemoteWrite; - let canonical = state_fingerprints(&request.queries[1], target).unwrap(); - let typed = &request.planner_candidate_forests[0]; - assert!(!canonical.is_empty()); - assert_eq!(canonical, state_fingerprints(&typed[1], target).unwrap()); - assert!(canonical.is_subset(&state_fingerprints(&typed[0], target).unwrap())); - } - - // A typed workload whose unretyped root reads a state of a retyped root - // is rejected rather than deploying two definitions of that state. - #[test] - fn typed_reselection_rejects_state_shared_with_unretyped_root() { - let (request, _) = - exact_workload_snapshot(&["avg_over_time(data[5m])", "sum_over_time(data[5m])"]) - .into_physical_compilation_request() - .unwrap(); - let inputs = planning_snapshot().physical_inputs; - let target = PhysicalDeploymentTarget::BackendLocalRemoteWrite; - let mut typed = request.planner_candidate_forests[0].clone(); - assert!(typed_reselection_forest( - &mut typed.clone(), - &request.queries, - &[true, true], - &needs_series_identity, - &inputs, - target, - ) - .is_ok()); - // The canonical `sum_over_time` stands in for a root that cannot be retyped. - typed[1] = request.queries[1].clone(); - let error = typed_reselection_forest( - &mut typed, - &request.queries, - &[true, false], - &needs_series_identity, - &inputs, - target, + let typed = state_fingerprints(&request.queries[1], target).unwrap(); + let mut canonical = request.queries.clone(); + select_logical_roots_with_scoped_evidence_and_trace( + &mut canonical, + request.canonical_roots.clone(), + &HashMap::new(), + &HashMap::new(), + &HashMap::new(), + None, + 1_000, ) - .unwrap_err(); - assert!(error.contains("shares state"), "{error}"); + .unwrap(); + assert!(!typed.is_empty()); + assert_eq!(typed, state_fingerprints(&canonical[1], target).unwrap()); + assert!(typed.is_subset(&state_fingerprints(&request.queries[0], target).unwrap())); } - // Reselection is gated on the missing series identity: a computation that - // fails to compile for another reason is not reselected. + // A sorted aggregation deploys from the first selected workload forest. #[test] - fn reselection_requires_the_series_identity_error() { - let query = "quantile_over_time(0.9,data[5m])/quantile_over_time(0.5,data[5m])"; - let mut snapshot = planning_snapshot(); - snapshot.query_workload.repeating_queries.as_mut().unwrap()[0].query = Query(query.into()); - let (request, _) = snapshot.into_physical_compilation_request().unwrap(); - assert!(request.planner_selection_trace.iter().all(|event| { - event["reselection"].is_null() - && event["stage"] != "planner.series_identity_reselection" - })); + fn sorted_aggregation_deploys_without_reselection() { + let query = "sort_desc(sum by (job) (rate(data[5m])))"; + let (request, environment) = exact_workload_snapshot(&[query]) + .into_physical_compilation_request() + .unwrap(); + let plan = DeploymentPlanCompiler + .compile_promql(request, environment) + .unwrap(); + let entry = plan.query_plan.lookup(query).unwrap(); + assert!(!entry.materialization_bindings().is_empty()); } // Quantile rank error does not certify relative error of a quotient. @@ -7292,13 +7168,13 @@ pub(crate) mod tests { ); assert!(!entry.materialization_bindings().is_empty()); } - // The grouped division is a Planner join over the two readouts. + // The grouped division combines the two readouts with Planner series matching. assert!(bundle .query_plan .entries .values() .flat_map(|entry| entry.nodes.values()) - .any(|node| !crate::query_plan::operator_parameters(node, "Join").is_empty())); + .any(|node| !crate::query_plan::operator_parameters(node, "SeriesBinary").is_empty())); } #[test] diff --git a/control_plane/src/physical/compiler/placement.rs b/control_plane/src/physical/compiler/placement.rs index 2bd7537ce..c2221fff3 100644 --- a/control_plane/src/physical/compiler/placement.rs +++ b/control_plane/src/physical/compiler/placement.rs @@ -61,11 +61,6 @@ impl Placement { .is_some_and(|mixed| mixed.retained.iter().any(|s| Rc::ptr_eq(s, state))) } - /// The root a mixed query compiles instead of Planner's selected root. - pub(super) fn root(&self, query: usize) -> Option<&Rc> { - self.mixed(query).map(|mixed| &mixed.root) - } - /// Native branches a mixed query maintains and reads as stored batches. pub(super) fn native_branches(&self, query: usize) -> &[Rc] { self.mixed(query).map_or(&[], |mixed| &mixed.branches) @@ -86,6 +81,7 @@ impl Placement { /// state's estimated bytes for every retained pane and partition at the /// summary-store price; an unknown size under a positive price leaves /// retention unpriced. +#[derive(Clone)] struct LifecycleCosts { costs: LifecycleUnitCosts, /// States the installed window layout retains. Without one, panes are @@ -763,11 +759,8 @@ fn raw_query_time_program(root: &QueryExpr) -> Result, rebuilt: Vec>, branches: Vec>, precompute: Vec, @@ -884,134 +877,40 @@ fn leaf_selector(state: &SummaryNode) -> Option<(Rc, QueryTimeOperato Some((Rc::clone(selector), scan)) } -/// A leaf state of an identity-typed root: its original and retyped nodes, -/// its raw leaf, and the `Scan` that reads that leaf at query time. +/// A state whose raw source already carries Planner's complete series identity. struct TypedState { - original: Rc, state: Rc, leaf: Rc, scan: QueryTimeOperator, } -/// `root` with each of `states`' raw selectors typed with the complete series -/// identity: raw rows read at query time carry it, and Planner's native -/// realization of a branch needs it. A per-series node keeps its input's -/// columns, so it gains the identity too; every other node is unchanged. -fn identity_typed( - root: &Rc, - states: &[Rc], -) -> Result<(Rc, Vec), String> { - use asap_physical_operators::physical_planner::promql_rows::SERIES_IDENTITY_COLUMN; - fn names(schema: &planner_types::post_asap::SummarySchema) -> Vec<&str> { - schema - .fields - .iter() - .map(|field| field.name.as_str()) - .collect() - } - // A readout or per-series state that passed its old input's columns - // through passes the identity on; any other node keeps its schema. - fn follow( - next: &mut SummaryNode, - old: &planner_types::post_asap::SummarySchema, - new: &planner_types::post_asap::SummarySchema, - ) { - let per_series = match &next.expr { - SummaryExpr::ValueOperation { .. } => true, - SummaryExpr::SummaryAgg { reduction, .. } => { - matches!(reduction, planner_types::pre_asap::Reduction::PerEntity) - } - _ => false, - }; - if per_series && names(&next.schema) == names(old) { - if let Some(identity) = new - .fields - .iter() - .find(|field| field.name == SERIES_IDENTITY_COLUMN) +/// Logical selection types all supported PromQL sources before placement. +/// Mixed placement only binds those typed leaves; it never rewrites their schemas. +fn typed_states(states: &[Rc]) -> Result, String> { + states + .iter() + .map(|state| { + let (selector, scan) = leaf_selector(state).ok_or("state has no raw selector")?; + if !selector + .output_schema() + .map_err(|error| error.to_string())? + .has_promql_series_identity() { - next.schema.fields.push(identity.clone()); + return Err("mixed input requires a logically typed series identity".into()); } - } - } - fn retype( - node: &Rc, - states: &[Rc], - typed: &mut Vec, - ) -> Result, String> { - let mut next = node.as_ref().clone(); - if states.iter().any(|state| Rc::ptr_eq(state, node)) { - let (selector, scan) = leaf_selector(node).ok_or("state has no raw selector")?; - let selector = - asap_physical_operators::physical_planner::promql_rows::with_series_identity( - &selector, - ) - .map_err(|e| e.to_string())?; - let SummaryExpr::SummaryAgg { child, .. } = &mut next.expr else { + let SummaryExpr::SummaryAgg { child, .. } = &state.expr else { unreachable!("leaf_selector matched a SummaryAgg") }; - let mut leaf = crate::planner_selection::keep_pre_asap(&selector) - .map_err(|e| e.to_string())? - .as_ref() - .clone(); - leaf.guarantee = child.guarantee.clone(); - let (old, leaf) = (Rc::clone(child), Rc::new(leaf)); - *child = Rc::clone(&leaf); - follow(&mut next, &old.schema, &leaf.schema); - let state = Rc::new(next); - typed.push(TypedState { - original: Rc::clone(node), - state: Rc::clone(&state), - leaf, + Ok(TypedState { + state: Rc::clone(state), + leaf: Rc::clone(child), scan, - }); - return Ok(state); - } - let unchanged = match &mut next.expr { - SummaryExpr::KeepPreAsap(_) => true, - SummaryExpr::ValueOperation { child, .. } | SummaryExpr::SummaryAgg { child, .. } => { - let old = Rc::clone(child); - *child = retype(child, states, typed)?; - let unchanged = Rc::ptr_eq(&old, child); - let new = Rc::clone(child); - follow(&mut next, &old.schema, &new.schema); - unchanged - } - // These nodes' schemas do not follow their inputs, so an input - // that gained the identity has no typed form here. - SummaryExpr::SummaryEstimate { summary_input, .. } => { - let old = Rc::clone(summary_input); - *summary_input = retype(summary_input, states, typed)?; - if old.schema != summary_input.schema { - return Err("an estimate over a per-series input cannot be retyped".into()); - } - Rc::ptr_eq(&old, summary_input) - } - SummaryExpr::BinaryOp { lhs, rhs, .. } => { - let (left, right) = (Rc::clone(lhs), Rc::clone(rhs)); - *lhs = retype(lhs, states, typed)?; - *rhs = retype(rhs, states, typed)?; - if left.schema != lhs.schema || right.schema != rhs.schema { - return Err("a per-series binary input cannot be retyped".into()); - } - Rc::ptr_eq(&left, lhs) && Rc::ptr_eq(&right, rhs) - } - _ => return Err("mixed placement supports no such operator".into()), - }; - Ok(if unchanged { - Rc::clone(node) - } else { - Rc::new(next) + }) }) - } - let mut typed = Vec::new(); - let root = retype(root, states, &mut typed)?; - if typed.len() != states.len() { - return Err("a state is not reachable as a leaf of the root".into()); - } - Ok((root, typed)) + .collect() } -/// Type `root` for a mixed plan that rebuilds `rebuilt` from raw rows and +/// Bind `root` for a mixed plan that rebuilds `rebuilt` from raw rows and /// reads each of `branches` as its stored native batch. Raw rows and Planner's /// native precompute of a branch both carry the complete series identity. fn mixed_placement( @@ -1023,12 +922,12 @@ fn mixed_placement( for branch in branches { leaves.extend(immutable_materialization_sources(branch).ok_or("branch has no sources")?); } - let (root, typed) = identity_typed(root, &leaves)?; + let typed = typed_states(&leaves)?; let (rebuilt, sources): (Vec<_>, Vec<_>) = typed .into_iter() - .partition(|state| rebuilt.iter().any(|s| Rc::ptr_eq(s, &state.original))); + .partition(|state| rebuilt.iter().any(|s| Rc::ptr_eq(s, &state.state))); // A branch is the Sum or sketch over readouts of its typed sources. - let typed_branches: Vec<_> = batch_branches(&root) + let typed_branches: Vec<_> = batch_branches(root) .into_iter() .filter(|branch| { immutable_materialization_sources(branch).is_some_and(|inputs| { @@ -1044,7 +943,7 @@ fn mixed_placement( } // Compilation numbers the same root identically when it installs the // plan, so these slots and precompute roots name its installed nodes. - let compiled = planner_types::post_asap::compile_post_asap_dag_with_node_ids(&root) + let compiled = planner_types::post_asap::compile_post_asap_dag_with_node_ids(root) .map_err(|e| e.to_string())?; let id = |node: &Rc| { compiled @@ -1075,7 +974,6 @@ fn mixed_placement( let mut retained = typed_branches.clone(); retained.extend(sources.into_iter().map(|state| state.state)); Ok(MixedPlacement { - root, rebuilt: rebuilt.into_iter().map(|state| state.state).collect(), branches: typed_branches, precompute, @@ -1258,6 +1156,41 @@ fn retime( Some(Rc::new(next)) } +/// Price the hypothetical retained native realization with the same window +/// preparation and selection used when that realization is installed. +fn native_retained_states( + root: &Rc, + state: &Rc, + query: &QueryCompilationInput, + environment: &PhysicalDeploymentContext, + inputs: &BackendLocalPhysicalInputs, +) -> Option { + let timing = planner_types::post_asap::ExecutionTiming::IngestionTime; + let retained_state = retime(state, state, timing)?; + let mut retained = query.clone(); + retained.selected_plan_root = retime(root, state, timing)?; + let physical = root_fixed_window_candidate(&retained.selected_plan_root).ok()?; + retained.retain(Some(physical)).ok()?; + prepare_window_implementations( + &mut retained, + &inputs.window_cost_model, + environment.target, + inputs.query_retention_margin_ms, + ) + .ok()?; + let states = collect_selected_materializations(&retained.selected_plan_root, true).ok()?; + let state = states.iter().find(|state| state.node == retained_state)?; + installed_window( + &retained, + state, + &states, + environment, + inputs.query_retention_margin_ms, + false, + ) + .map(|(_, count)| count) +} + /// A Planner candidate with a native physical realization after its /// readout-built states are placed by lifecycle. pub(super) struct TimedCandidate { @@ -1278,6 +1211,7 @@ pub(super) fn time_native_candidate( data: &DataWorkload, index: usize, environment: &PhysicalDeploymentContext, + inputs: &BackendLocalPhysicalInputs, ) -> Option { use planner_types::post_asap::ExecutionTiming; let untimed = || { @@ -1288,8 +1222,8 @@ pub(super) fn time_native_candidate( }) }; let lifecycle = &query.summary_lifecycle_inputs; - // Windows are prepared from the timed root, so the installed layout is not - // known while its timing is being chosen. + // Enumerate lifecycle choices first; each retained native alternative is + // repriced below using its hypothetical installed window layout. let model = LifecycleCosts { costs: lifecycle.costs.clone(), retained_states: None, @@ -1322,11 +1256,32 @@ pub(super) fn time_native_candidate( let mut trace = Vec::new(); let mut maintained_over_readouts = false; for deployment in candidates.deployments() { - let retained = alternative_cost( - deployment, - &SummaryMaintenanceLifecycle::ContinuouslyMaintained, - ); let lifecycle = if over_readouts(&deployment.summary) { + let retained_states = + native_retained_states(root, &deployment.summary, query, environment, inputs); + let mut retained_model = model.clone(); + retained_model.retained_states = retained_states; + let retained = retained_states.and_then(|_| { + let priced = enumerate_summary_maintenance_lifecycles( + Rc::clone(root), + WorkloadDemand::new_with_data(workload, data, std::slice::from_ref(&index)), + environment.observed_at_unix_ms, + Some(Horizon(lifecycle.horizon_seconds)), + SummaryMaintenanceLifecycleCapabilities { + supports_ephemeral: true, + supports_prepared: false, + supports_shared: false, + supports_continuously_maintained: true, + }, + &retained_model, + ) + .ok()?; + let priced = priced + .deployments() + .iter() + .find(|priced| priced.summary == deployment.summary)?; + alternative_cost(priced, &SummaryMaintenanceLifecycle::ContinuouslyMaintained) + }); let rebuilt = alternative_cost(deployment, &SummaryMaintenanceLifecycle::Ephemeral); let ephemeral = rebuild_is_cheaper(retained, rebuilt); maintained_over_readouts |= !ephemeral; @@ -1340,6 +1295,7 @@ pub(super) fn time_native_candidate( // placement of the states it installs as `lifecycle_placement`. trace.push(json!({ "stage": "deployment.native_candidate_placement", + "retained_states": retained_states, "query_ids": [&query.query_id], "candidate_root_id": crate::planner_selection::explained_root_id(root, &query.accuracy_target), "logical_root_id": crate::planner_selection::explained_root_id(&deployment.summary, &query.accuracy_target), diff --git a/control_plane/src/planner_selection.rs b/control_plane/src/planner_selection.rs index 8bcc79fdf..ec7961ff9 100644 --- a/control_plane/src/planner_selection.rs +++ b/control_plane/src/planner_selection.rs @@ -670,20 +670,42 @@ mod workload_tests { assert!(saw_unknown); } - // JSON must not alias NaN and infinity through its null representation. + // Explicit infinity encodings retain distinct stable explanation identities. #[test] - fn explain_nonfinite_identity_is_unavailable() { - for value in [f64::NAN, f64::INFINITY, f64::NEG_INFINITY] { + fn explain_infinity_identities_are_distinct_and_stable() { + let mut identities = std::collections::BTreeSet::new(); + for value in [f64::INFINITY, f64::NEG_INFINITY] { let root = QueryExpr::Literal(planner_types::pre_asap::ScalarValue::Float64(value)); - assert!(replacement_identity( + let identity = replacement_identity( &root, &Replacement::Rewrite(Rc::new(root.clone())), - &AccuracyTarget::Exact + &AccuracyTarget::Exact, ) - .is_none()); + .unwrap(); + assert_eq!( + Some(identity.clone()), + replacement_identity( + &root, + &Replacement::Rewrite(Rc::new(root.clone())), + &AccuracyTarget::Exact + ) + ); + assert!(identities.insert(identity)); } } + // NaN is not equal to itself, so identity validation remains conservative. + #[test] + fn explain_nan_identity_is_unavailable() { + let root = QueryExpr::Literal(planner_types::pre_asap::ScalarValue::Float64(f64::NAN)); + assert!(replacement_identity( + &root, + &Replacement::Rewrite(Rc::new(root.clone())), + &AccuracyTarget::Exact, + ) + .is_none()); + } + // Hashing excludes incidental allocation sharing but retains operand roles. #[test] fn explain_summary_identity_preserves_roles_and_ignores_rc_sharing() { diff --git a/control_plane/src/query_plan.rs b/control_plane/src/query_plan.rs index faf45a98a..5cc2e0b34 100644 --- a/control_plane/src/query_plan.rs +++ b/control_plane/src/query_plan.rs @@ -1077,11 +1077,9 @@ mod catalog_binding_tests { mod tests { use super::*; - // Planner's guarded average divides two per-series readouts; until - // Planner matches per-series rows, lowering refuses it instead of - // computing it in the backend. + // Planner compiles guarded per-series average division into a physical fragment. #[test] - fn per_series_guarded_division_is_not_lowered_locally() { + fn per_series_guarded_division_uses_planner_fragment() { let query = "avg_over_time(m[5m])"; let canonical = crate::query_parser::parse_query_expr_canonical( query, @@ -1093,7 +1091,7 @@ mod tests { panic!("expected the Planner's average rewrite"); }; assert!(operator.checked_finite_division); - let error = compile_bound_mapped( + let entry = compile_bound_mapped( "guarded".into(), query.into(), &root, @@ -1119,11 +1117,14 @@ mod tests { }, |_, _| {}, ) - .unwrap_err(); - assert!( - matches!(error, QueryPlanError::UnsupportedNode(_)), - "{error}" - ); + .unwrap(); + let QueryPlanNode::PhysicalFragment { dag, .. } = &entry.nodes[&entry.root] else { + panic!("expected compiled Planner fragment") + }; + let document: serde_json::Value = serde_json::from_slice(dag).unwrap(); + assert!(document + .to_string() + .contains("\"checked_finite_division\":true")); } // Every o11y corpus query compiles as backend readouts under Planner @@ -1164,7 +1165,9 @@ mod tests { for entry in plan.query_plan.entries.values() { for node in entry.nodes.values() { match node { - QueryPlanNode::PhysicalFragment { .. } => fragments += 1, + QueryPlanNode::PhysicalFragment { .. } | QueryPlanNode::Physical { .. } => { + fragments += 1 + } QueryPlanNode::ReadMaterialization { .. } | QueryPlanNode::ExactReadout { .. } | QueryPlanNode::SummaryEstimate { .. } diff --git a/crates/asap_summary_state/src/lib.rs b/crates/asap_summary_state/src/lib.rs index 60e394581..4fc85da0a 100644 --- a/crates/asap_summary_state/src/lib.rs +++ b/crates/asap_summary_state/src/lib.rs @@ -10,7 +10,6 @@ pub use aggregation_type::AggregationType; pub mod codec; pub mod stored_state; -pub mod univmon; pub use asap_physical_operators::{AggregateCore, KeyByLabelValues, Measurement, Statistic}; pub use stored_state::codec::StoredState; diff --git a/crates/asap_summary_state/src/stored_state/codec.rs b/crates/asap_summary_state/src/stored_state/codec.rs index 6f94f5932..8ac81460e 100644 --- a/crates/asap_summary_state/src/stored_state/codec.rs +++ b/crates/asap_summary_state/src/stored_state/codec.rs @@ -4,9 +4,9 @@ //! only place that names their persisted type tags and byte encodings; Planner //! owns their update, merge and estimate. use super::native::NativeSummaryOutput; -use crate::univmon::UnivMonAccumulator; -use crate::{AggregationType, KeyByLabelValues}; +use crate::AggregationType; use asap_physical_operators::summary_kernels as k; +use asap_physical_operators::summary_kernels::univmon::UnivMonAccumulator; use asap_physical_operators::summary_kernels::weighted_frequency::{ FrequencyAlgorithm, WeightedFrequency, }; @@ -18,11 +18,12 @@ use std::collections::HashMap; pub type Error = Box; /// Persisted type tag of Planner's exact state (its named msgpack form). -pub const EXACT_V1: &str = "PlannerExactAccumulatorV1"; +pub const EXACT_V2: &str = "PlannerExactAccumulatorV2"; /// Exact-state tags written before the store kept Planner kernels. Their byte /// layouts are no longer decoded. const RETIRED_EXACT_TAGS: &[&str] = &[ + "PlannerExactAccumulatorV1", "SumAccumulator", "IncreaseAccumulator", "MinAccumulator", @@ -80,20 +81,8 @@ fn view(state: &dyn AggregateCore) -> Option> { None } -/// Planner's weighted frequency state is serde-transparent over the sketchlib -/// kernel, which owns the persisted `WeightedFrequencyV1` bytes. -// Planner keeps the kernel and its algorithm private; this round trip is the -// only public access until it exposes them. -pub(crate) fn frequency_kernel( - state: &WeightedFrequency, -) -> Result { - Ok(rmp_serde::from_slice(&rmp_serde::to_vec(state)?)?) -} - pub(crate) fn frequency_state(bytes: &[u8]) -> Result { - let kernel = asap_sketchlib::WeightedFrequency::from_bytes(bytes) - .map_err(|e| format!("deserialize weighted frequency: {e:?}"))?; - Ok(rmp_serde::from_slice(&rmp_serde::to_vec(&kernel)?)?) + Ok(WeightedFrequency::from_bytes(bytes)?) } fn exact_aggregation_type(state: &k::exact::ExactAccumulator) -> AggregationType { @@ -126,12 +115,12 @@ pub trait StoredState { impl<'a> StoredState for dyn AggregateCore + 'a { fn type_name(&self) -> &'static str { match view(self) { - Some(View::Exact(_)) => EXACT_V1, - Some(View::Dd(_)) => "DDSketchAccumulator", - Some(View::Hll(_)) => "HllSketchAccumulator", + Some(View::Exact(_)) => EXACT_V2, + Some(View::Dd(_)) => "DDSketchAccumulatorV2", + Some(View::Hll(_)) => "HllSketchAccumulatorV2", Some(View::Kll(_)) => "DatasketchesKLLAccumulator", - Some(View::Cms(_)) => "CountMinSketchAccumulator", - Some(View::Cs(_)) => "CountSketchAccumulator", + Some(View::Cms(_)) => "CountMinSketchAccumulatorV2", + Some(View::Cs(_)) => "CountSketchAccumulatorV2", Some(View::CmsHeap(_)) => "CountMinSketchWithHeapAccumulator", Some(View::CsHeap(_)) => "CountSketchWithHeapAccumulator", Some(View::Hydra(_)) => "HydraKllSketchAccumulator", @@ -154,8 +143,8 @@ impl<'a> StoredState for dyn AggregateCore + 'a { Some(View::CmsHeap(_)) => T::CountMinSketchWithHeap, Some(View::CsHeap(_)) => T::CountSketchWithHeap, Some(View::Hydra(_)) => T::HydraKLL, - Some(View::Frequency(state)) => match frequency_kernel(state).map(|k| k.algorithm()) { - Ok(FrequencyAlgorithm::CountSketch) => T::CountSketchWithHeap, + Some(View::Frequency(state)) => match state.algorithm() { + FrequencyAlgorithm::CountSketch => T::CountSketchWithHeap, _ => T::CountMinSketchWithHeap, }, Some(View::UnivMon(_)) => T::UnivMon, @@ -168,16 +157,24 @@ impl<'a> StoredState for dyn AggregateCore + 'a { Ok( match view(self).ok_or("Planner state has no stored codec")? { View::Exact(state) => rmp_serde::to_vec_named(state)?, - View::Dd(state) => state.inner.to_msgpack()?, - View::Hll(state) => state.inner.to_msgpack()?, + View::Dd(state) => { + rmp_serde::to_vec(&(state.sample_p(), state.inner.to_msgpack()?))? + } + View::Hll(state) => { + rmp_serde::to_vec(&(state.sample_p(), state.inner.to_msgpack()?))? + } View::Kll(state) => state.inner.to_msgpack()?, - View::Cms(state) => state.inner.to_msgpack()?, - View::Cs(state) => state.inner.to_msgpack()?, + View::Cms(state) => { + rmp_serde::to_vec(&(state.sample_p(), state.inner.to_msgpack()?))? + } + View::Cs(state) => { + rmp_serde::to_vec(&(state.sample_p(), state.inner.to_msgpack()?))? + } View::CmsHeap(state) => state.inner.to_msgpack()?, View::CsHeap(state) => state.inner.to_msgpack()?, View::Hydra(state) => state.inner.to_msgpack()?, - View::Frequency(state) => frequency_kernel(state)?.to_bytes(), - View::UnivMon(state) => state.to_bytes()?, + View::Frequency(state) => state.to_bytes(), + View::UnivMon(state) => state.sketch().serialize_to_bytes()?, View::Native(state) => state.bytes().to_vec(), }, ) @@ -214,7 +211,7 @@ pub fn range_ms(parameters: &HashMap) -> Result Result<(), Error> { view(state) .map(|_| ()) @@ -226,22 +223,38 @@ pub fn check_storable(state: &dyn AggregateCore) -> Result<(), Error> { pub fn decode(type_name: &str, bytes: &[u8]) -> Result, Error> { use super::decoders as d; Ok(match type_name { - EXACT_V1 => Box::new(decode_exact(bytes)?), - "DDSketchAccumulator" => Box::new(k::DDSketchAccumulator { - inner: d::ddsketch_from_msgpack(bytes)?, - }), - "HllSketchAccumulator" => Box::new(k::HllSketchAccumulator { - inner: d::hll_from_msgpack(bytes)?, - }), + EXACT_V2 => Box::new(decode_exact(bytes)?), + "DDSketchAccumulatorV2" => { + let (sample_p, sketch): (f64, Vec) = rmp_serde::from_slice(bytes)?; + Box::new(k::DDSketchAccumulator::from_sketch( + d::ddsketch_from_msgpack(&sketch)?, + sample_p, + )?) + } + "HllSketchAccumulatorV2" => { + let (sample_p, sketch): (f64, Vec) = rmp_serde::from_slice(bytes)?; + Box::new(k::HllSketchAccumulator::from_sketch( + d::hll_from_msgpack(&sketch)?, + sample_p, + )?) + } "DatasketchesKLLAccumulator" => Box::new(k::DatasketchesKLLAccumulator { inner: d::kll_from_msgpack(bytes)?, }), - "CountMinSketchAccumulator" => Box::new(k::CountMinSketchAccumulator { - inner: d::cms_from_msgpack(bytes)?, - }), - "CountSketchAccumulator" => Box::new(k::CountSketchAccumulator { - inner: d::cs_from_msgpack(bytes)?, - }), + "CountMinSketchAccumulatorV2" => { + let (sample_p, sketch): (f64, Vec) = rmp_serde::from_slice(bytes)?; + Box::new(k::CountMinSketchAccumulator::from_sketch( + d::cms_from_msgpack(&sketch)?, + sample_p, + )?) + } + "CountSketchAccumulatorV2" => { + let (sample_p, sketch): (f64, Vec) = rmp_serde::from_slice(bytes)?; + Box::new(k::CountSketchAccumulator::from_sketch( + d::cs_from_msgpack(&sketch)?, + sample_p, + )?) + } "CountMinSketchWithHeapAccumulator" => Box::new(k::CountMinSketchWithHeapAccumulator { inner: d::cms_with_heap_from_msgpack(bytes)?, }), @@ -252,11 +265,22 @@ pub fn decode(type_name: &str, bytes: &[u8]) -> Result, E inner: asap_sketchlib::HydraKllSketch::from_msgpack(bytes)?, }), "WeightedFrequency" => Box::new(frequency_state(bytes)?), - "UnivMonAccumulator" => Box::new(UnivMonAccumulator::from_bytes(bytes)?), + "UnivMonAccumulator" => Box::new(UnivMonAccumulator::from_sketch( + asap_sketchlib::UnivMon::deserialize_from_bytes(bytes)?, + )?), + "DDSketchAccumulator" + | "HllSketchAccumulator" + | "CountMinSketchAccumulator" + | "CountSketchAccumulator" => { + return Err(format!( + "stored format {type_name} is retired; sampled sketch storage requires V2" + ) + .into()); + } retired if is_retired_exact(retired) => { return Err(format!( "stored format {retired} is retired and no longer decoded; \ - exact state is stored as {EXACT_V1}" + exact state is stored as {EXACT_V2}" ) .into()) } @@ -276,18 +300,34 @@ pub fn decode_envelope(bytes: &[u8]) -> Result, Error> { Some(SketchState::Kll(_)) => Box::new(k::DatasketchesKLLAccumulator { inner: d::kll_from_proto(bytes)?, }), - Some(SketchState::Ddsketch(_)) => Box::new(k::DDSketchAccumulator { - inner: d::ddsketch_from_proto(bytes)?, - }), - Some(SketchState::Hll(_)) => Box::new(k::HllSketchAccumulator { - inner: d::hll_from_proto(bytes)?, - }), - Some(SketchState::CountMin(_)) => Box::new(k::CountMinSketchAccumulator { - inner: d::cms_from_proto(bytes)?, - }), - Some(SketchState::CountSketch(_)) => Box::new(k::CountSketchAccumulator { - inner: d::cs_from_proto(bytes)?, - }), + Some(SketchState::Ddsketch(_)) => Box::new( + k::DDSketchAccumulator::from_sketch( + d::ddsketch_from_proto(bytes)?, + d::sample_probability(bytes)?, + ) + .expect("validated sampling probability"), + ), + Some(SketchState::Hll(_)) => Box::new( + k::HllSketchAccumulator::from_sketch( + d::hll_from_proto(bytes)?, + d::sample_probability(bytes)?, + ) + .expect("validated sampling probability"), + ), + Some(SketchState::CountMin(_)) => Box::new( + k::CountMinSketchAccumulator::from_sketch( + d::cms_from_proto(bytes)?, + d::sample_probability(bytes)?, + ) + .expect("validated sampling probability"), + ), + Some(SketchState::CountSketch(_)) => Box::new( + k::CountSketchAccumulator::from_sketch( + d::cs_from_proto(bytes)?, + d::sample_probability(bytes)?, + ) + .expect("validated sampling probability"), + ), Some(other) => { let family = match other { SketchState::Univmon(_) => "UnivMon", @@ -305,44 +345,6 @@ pub fn decode_envelope(bytes: &[u8]) -> Result, Error> { /// Decode Planner's exact state, rejecting a payload whose population states /// differ from its declared family. pub fn decode_exact(bytes: &[u8]) -> Result { - // Planner's exact state is serde-derived without validation; this mirror - // of its persisted shape checks each population against the family. - #[derive(serde::Deserialize)] - struct Shape { - family: SummaryFamilyType, - scalar: Scalar, - keyed: Option>, - } - #[derive(serde::Deserialize)] - enum Scalar { - Sum(serde::de::IgnoredAny), - Count(serde::de::IgnoredAny), - Min(serde::de::IgnoredAny), - Max(serde::de::IgnoredAny), - Counter(serde::de::IgnoredAny), - } - let shape: Shape = rmp_serde::from_slice(bytes)?; - // Planner accepts only matching (kind, params) exact families. - k::exact::ExactAccumulator::new(shape.family.clone(), shape.keyed.is_some())?; - let SummaryFamilyType::ExactAggregate(expected, _) = &shape.family else { - return Err(format!("{:?} is not an exact family", shape.family).into()); - }; - let matches = |scalar: &Scalar| match scalar { - Scalar::Sum(_) => *expected == ExactKind::Sum, - Scalar::Count(_) => *expected == ExactKind::Count, - Scalar::Min(_) => *expected == ExactKind::Min, - Scalar::Max(_) => *expected == ExactKind::Max, - Scalar::Counter(_) => matches!(expected, ExactKind::Rate | ExactKind::Increase), - }; - if !matches(&shape.scalar) - || shape - .keyed - .iter() - .flat_map(HashMap::values) - .any(|s| !matches(s)) - { - return Err("exact payload differs from declared Planner family".into()); - } Ok(rmp_serde::from_slice(bytes)?) } @@ -355,18 +357,31 @@ pub fn empty_like(state: &dyn AggregateCore) -> Result, E }; Ok( match view(state).ok_or("Planner state has no stored codec")? { - View::Dd(s) => Box::new(k::DDSketchAccumulator { - inner: DdSketch::new(s.inner.alpha), - }), - View::Hll(s) => Box::new(k::HllSketchAccumulator { - inner: HllSketch::new(s.inner.variant, s.inner.precision), - }), - View::Cms(s) => Box::new(k::CountMinSketchAccumulator { - inner: CountMinSketch::new(s.inner.rows(), s.inner.cols()), - }), - View::Cs(s) => Box::new(k::CountSketchAccumulator { - inner: CountSketch::new(s.inner.rows, s.inner.cols), - }), + View::Dd(s) => Box::new( + k::DDSketchAccumulator::from_sketch(DdSketch::new(s.inner.alpha), s.sample_p()) + .expect("validated sampling probability"), + ), + View::Hll(s) => Box::new( + k::HllSketchAccumulator::from_sketch( + HllSketch::new(s.inner.variant, s.inner.precision), + s.sample_p(), + ) + .expect("validated sampling probability"), + ), + View::Cms(s) => Box::new( + k::CountMinSketchAccumulator::from_sketch( + CountMinSketch::new(s.inner.rows(), s.inner.cols()), + s.sample_p(), + ) + .expect("validated sampling probability"), + ), + View::Cs(s) => Box::new( + k::CountSketchAccumulator::from_sketch( + CountSketch::new(s.inner.rows, s.inner.cols), + s.sample_p(), + ) + .expect("validated sampling probability"), + ), View::CmsHeap(s) => Box::new(k::CountMinSketchWithHeapAccumulator { inner: CountMinSketchWithHeap::new( s.inner.rows(), @@ -378,9 +393,8 @@ pub fn empty_like(state: &dyn AggregateCore) -> Result, E inner: CountSketchWithHeap::new(s.inner.rows(), s.inner.cols(), s.inner.heap_size), }), View::UnivMon(s) => { - let mut empty = s.clone(); - empty.clear(); - Box::new(empty) + let (heap, rows, cols, layers) = s.dimensions(); + Box::new(UnivMonAccumulator::new(heap, rows, cols, layers)?) } _ => return Err("only delta-capable sketch families reset to empty".into()), }, @@ -390,28 +404,31 @@ pub fn empty_like(state: &dyn AggregateCore) -> Result, E #[cfg(test)] mod tests { use super::*; - use crate::Statistic; + use crate::{KeyByLabelValues, Statistic}; use asap_physical_operators::values::Value; use planner_types::{post_asap::SketchQuery, pre_asap::ColumnRef}; - /// Stored bytes written by the pre-Planner-kernel backend, one per family. + /// Golden payloads for each supported stored family; sampled codecs wrap sketch bytes. const GOLDEN: &[(&str, &str)] = &[ - ("DDSketchAccumulator", "93cb3f847ae147ae147bdc01040000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000001000000000000000000000000000000000000000000000000000000000000000000020000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000001d0c0"), - ("HllSketchAccumulator", "415341507631010201010000013a0000001388b06d657461646174615f76657273696f6e01af686173685f70726f66696c655f6964bc70726f6a656374617361702e787868332e736565646c6973742e7631ae686173685f616c676f726974686dab787868335f36345f313238af736565645f64657269766174696f6eb4736565645f6c6973745f696e6465785f77726170ae696e7075745f656e636f64696e67b470726f6a656374617361702e696e7075742e7631a9736565645f6c697374dc0014cecafe3553cf000000ade3415118ce8cc70208ce2f024b2bce451a3df5ce6a09e667cebb67ae85ce3c6ef372cea54ff53ace510e527fce9b05688cce1f83d9abce5be0cd19cecbbb9d5dce629a292ace9159015ace152fecd8ce67332667ce8eb44a87cedb0c2e0db463616e6f6e6963616c5f736565645f696e64657805a9707265636973696f6e0491c41000000300010101000000010000000000"), + ("DDSketchAccumulatorV2", "93cb3f847ae147ae147bdc01040000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000001000000000000000000000000000000000000000000000000000000000000000000020000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000001d0c0"), + ("HllSketchAccumulatorV2", "415341507631010201010000013a0000001388b06d657461646174615f76657273696f6e01af686173685f70726f66696c655f6964bc70726f6a656374617361702e787868332e736565646c6973742e7631ae686173685f616c676f726974686dab787868335f36345f313238af736565645f64657269766174696f6eb4736565645f6c6973745f696e6465785f77726170ae696e7075745f656e636f64696e67b470726f6a656374617361702e696e7075742e7631a9736565645f6c697374dc0014cecafe3553cf000000ade3415118ce8cc70208ce2f024b2bce451a3df5ce6a09e667cebb67ae85ce3c6ef372cea54ff53ace510e527fce9b05688cce1f83d9abce5be0cd19cecbbb9d5dce629a292ace9159015ace152fecd8ce67332667ce8eb44a87cedb0c2e0db463616e6f6e6963616c5f736565645f696e64657805a9707265636973696f6e0491c41000000300010101000000010000000000"), ("DatasketchesKLLAccumulator", "92ccc8dc006641534150763101020600000000280000002ccc84ccb06d657461646174615f76657273696f6e01cca16bccccccc8cca16d08cca96974656d5f74797065cca3663634cc93cc920003cc93cccb4008000000000000cccb3fccf0000000000000cccb4000000000000000cc93cccf560f2acc9b7e7e3ccca80000"), - ("CountMinSketchAccumulator", "939298cb0000000000000000cb4000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb000000000000000098cb0000000000000000cb4000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb00000000000000000208"), - ("CountSketchAccumulator", "9403089398cb0000000000000000cb4000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb000000000000000098cb0000000000000000cbc000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb000000000000000098cb0000000000000000cb4000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb000000000000000090"), + ("CountMinSketchAccumulatorV2", "939298cb0000000000000000cb4000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb000000000000000098cb0000000000000000cb4000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb00000000000000000208"), + ("CountSketchAccumulatorV2", "9403089398cb0000000000000000cb4000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb000000000000000098cb0000000000000000cbc000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb000000000000000098cb0000000000000000cb4000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb000000000000000090"), ("CountMinSketchWithHeapAccumulator", "93939298cb0000000000000000cb4008000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb000000000000000098cb4008000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb000000000000000002089192a161cb400800000000000002"), ("CountSketchWithHeapAccumulator", "93939398cb0000000000000000cb4008000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb000000000000000098cbc008000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb000000000000000098cb0000000000000000cb0000000000000000cbc008000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb0000000000000000cb000000000000000003089192a161cb400800000000000002"), - (EXACT_V1, "83a666616d696c7981ae457861637441676772656761746592a353756da353756da67363616c617281a353756dcb4012000000000000a56b65796564c0"), - (EXACT_V1, "83a666616d696c7981ae457861637441676772656761746592a5436f756e74a5436f756e74a67363616c617281a5436f756e7400a56b657965648181a66c6162656c7391a16181a5436f756e7402"), - (EXACT_V1, "83a666616d696c7981ae457861637441676772656761746592a452617465a452617465a67363616c617281a7436f756e74657286b47374617274696e675f6d6561737572656d656e7481a576616c7565cb4024000000000000b27374617274696e675f74696d657374616d70cd03e8b56c6173745f7365656e5f6d6561737572656d656e7481a576616c7565cb4039000000000000b36c6173745f7365656e5f74696d657374616d70cd07d0ae746f74616c5f696e637265617365cb402e000000000000ac73616d706c655f636f756e7402a56b65796564c0"), + (EXACT_V2, "83a666616d696c7981ae457861637441676772656761746592a353756da353756da67363616c617281a553756d563282a373756dcb4012000000000000ac636f6d70656e736174696f6ecb0000000000000000a56b65796564c0"), + (EXACT_V2, "83a666616d696c7981ae457861637441676772656761746592a5436f756e74a5436f756e74a67363616c617281a5436f756e7400a56b657965648181a66c6162656c7391a16181a5436f756e7402"), + (EXACT_V2, "83a666616d696c7981ae457861637441676772656761746592a452617465a452617465a67363616c617281a7436f756e74657286b47374617274696e675f6d6561737572656d656e7481a576616c7565cb4024000000000000b27374617274696e675f74696d657374616d70cd03e8b56c6173745f7365656e5f6d6561737572656d656e7481a576616c7565cb4039000000000000b36c6173745f7365656e5f74696d657374616d70cd07d0ae746f74616c5f696e637265617365cb402e000000000000ac73616d706c655f636f756e7402a56b65796564c0"), ("UnivMonAccumulator", "4153415076310102100000000155000000898bb06d657461646174615f76657273696f6e01af686173685f70726f66696c655f6964bc70726f6a656374617361702e787868332e736565646c6973742e7631ae686173685f616c676f726974686dab787868335f36345f313238af736565645f64657269766174696f6eb4736565645f6c6973745f696e6465785f77726170ae696e7075745f656e636f64696e67b470726f6a656374617361702e696e7075742e7631a9736565645f6c697374dc0014cecafe3553cf000000ade3415118ce8cc70208ce2f024b2bce451a3df5ce6a09e667cebb67ae85ce3c6ef372cea54ff53ace510e527fce9b05688cce1f83d9abce5be0cd19cecbbb9d5dce629a292ace9159015ace152fecd8ce67332667ce8eb44a87cedb0c2e0daa6c617965725f73697a6502aa736b657463685f726f7703aa736b657463685f636f6c10a9686561705f73697a6504a86b65795f74797065a375363498dc0060000000000000ff00000000000000010000000000000001000000000000ff00000000000000000000000000000001ff000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000009602020200000092020092cf3ff0000000000000cf400000000000000092010192c3c30201"), ("WeightedFrequency", "415341502d57465245512d31000000000008000000000000000200000000000000020000000000000010000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000e03f0000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000e03f01000000000000000100000000000000040000000100000000000000611500000000000000010000000000000004000000010000000000000061000000000000e03f"), ]; fn golden(tag: &str, hex_bytes: &str) -> (Box, Vec) { - let bytes = hex::decode(hex_bytes).unwrap(); + let mut bytes = hex::decode(hex_bytes).unwrap(); + if tag.ends_with("V2") && tag != EXACT_V2 { + bytes = rmp_serde::to_vec(&(1.0_f64, bytes)).unwrap(); + } (decode(tag, &bytes).unwrap(), bytes) } fn item() -> Option { @@ -519,13 +536,100 @@ mod tests { assert!(decode_exact(&bytes).is_err()); } - // A state without a stored codec cannot enter the store. + // Sampled edge counts retain their probability through durable storage. + #[test] + fn sampled_envelope_storage_preserves_scaled_count() { + use asap_sketchlib::proto::sketchlib::SketchEnvelope; + use prost::Message; + let mut sketch = asap_sketchlib::DdSketch::new(0.01); + sketch.update(3.0); + let mut envelope = + SketchEnvelope::decode(asap_sketch_codec::encode_ddsketch(&sketch).as_slice()).unwrap(); + for (p, expected) in [(0.0, 1.0), (1.0, 1.0), (0.25, 4.0)] { + envelope.sample_p = p; + let state = decode_envelope(&envelope.encode_to_vec()).unwrap(); + let restored = decode(state.type_name(), &state.encode().unwrap()).unwrap(); + assert_eq!(read(restored.as_ref(), Statistic::Count, None), expected); + } + for p in [-1.0, 1.1, f64::NAN, f64::INFINITY] { + envelope.sample_p = p; + assert!(decode_envelope(&envelope.encode_to_vec()).is_err()); + } + assert!(decode("DDSketchAccumulator", &sketch.to_msgpack().unwrap()) + .err() + .unwrap() + .to_string() + .contains("retired")); + } + + // Every sampled frequency/cardinality family retains scaled readouts in storage. + #[test] + fn sampled_frequency_and_hll_storage_keep_probability() { + let key = KeyByLabelValues::new_with_labels(vec!["a".into()]); + let mut cms = asap_sketchlib::CountMinSketch::new(2, 16); + cms.update("a", 3.0); + let mut cs = asap_sketchlib::CountSketch::new(3, 16); + cs.update("a", 3.0); + let states: Vec> = vec![ + Box::new(k::CountMinSketchAccumulator::from_sketch(cms, 0.25).unwrap()), + Box::new(k::CountSketchAccumulator::from_sketch(cs, 0.25).unwrap()), + ]; + for state in states { + let restored = decode(state.type_name(), &state.encode().unwrap()).unwrap(); + assert_eq!( + read(restored.as_ref(), Statistic::Count, Some(key.clone())), + 12.0 + ); + } + let state = k::HllSketchAccumulator::from_sketch( + asap_sketchlib::HllSketch::new(asap_sketchlib::HllVariant::Regular, 4), + 0.25, + ) + .unwrap(); + let state = &state as &dyn AggregateCore; + let restored = decode(state.type_name(), &state.encode().unwrap()).unwrap(); + assert_eq!( + restored + .as_any() + .downcast_ref::() + .unwrap() + .sample_p(), + 0.25 + ); + } + + // Window resets discard counts while retaining a series' sampling probability. + #[test] + fn sampled_reset_keeps_probability() { + let state = + k::DDSketchAccumulator::from_sketch(asap_sketchlib::DdSketch::new(0.01), 0.25).unwrap(); + let empty = empty_like(&state).unwrap(); + assert_eq!( + empty + .as_any() + .downcast_ref::() + .unwrap() + .sample_p(), + 0.25 + ); + } + + // Persisted terminal-mode UnivMon sketches cannot enter the standard-update kernel. + #[test] + fn terminal_univmon_storage_is_rejected() { + let mut sketch = asap_sketchlib::UnivMon::init_univmon(4, 3, 16, 2); + sketch.fast_insert(&asap_sketchlib::DataInput::U64(1), 1); + assert!(decode("UnivMonAccumulator", &sketch.serialize_to_bytes().unwrap()).is_err()); + } + + // Planner UnivMon states persist directly without a backend kernel shim. #[test] - fn planner_univmon_has_no_stored_codec() { - let planner = k::univmon::UnivMonAccumulator::new(4, 3, 16, 2).unwrap(); - assert!(check_storable(&planner).is_err()); - assert!((&planner as &dyn AggregateCore).encode().is_err()); - assert!(check_storable(&UnivMonAccumulator::new(4, 3, 16, 2).unwrap()).is_ok()); + fn planner_univmon_round_trips_through_storage() { + let state = k::univmon::UnivMonAccumulator::new(4, 3, 16, 2).unwrap(); + check_storable(&state).unwrap(); + let bytes = (&state as &dyn AggregateCore).encode().unwrap(); + let restored = decode("UnivMonAccumulator", &bytes).unwrap(); + assert!(restored.as_any().is::()); } // An OTLP envelope attribute decodes into its family's Planner kernel; diff --git a/crates/asap_summary_state/src/stored_state/decoders.rs b/crates/asap_summary_state/src/stored_state/decoders.rs index 53de7893c..a1883d75b 100644 --- a/crates/asap_summary_state/src/stored_state/decoders.rs +++ b/crates/asap_summary_state/src/stored_state/decoders.rs @@ -11,17 +11,23 @@ use asap_sketchlib::{ }; use prost::Message; -/// Planner kernels carry no edge sampling probability, so a sampled frame -/// (`0 < sample_p < 1`) would silently read as unscaled counts. `0` (proto3 -/// default) and `1` both mean unsampled. -fn unsampled(what: &str, sample_p: f64) -> Result<(), String> { - if sample_p.is_finite() && sample_p > 0.0 && sample_p < 1.0 { +/// Validate edge probability; proto3's omitted zero means unsampled. +fn checked_probability(sample_p: f64) -> Result { + let p = if sample_p == 0.0 { 1.0 } else { sample_p }; + if !p.is_finite() || p <= 0.0 || p > 1.0 { return Err(format!( - "{what} frame is edge-sampled (sample_p={sample_p}); sampled sketch \ - state is not supported by Planner kernels" + "invalid sketch sample_p={sample_p}; expected 0 < p <= 1" )); } - Ok(()) + Ok(p) +} + +/// Sampling metadata of a full envelope; bare state/delta frames are unsampled. +pub fn sample_probability(buffer: &[u8]) -> Result { + match SketchEnvelope::decode(buffer) { + Ok(envelope) if envelope.sketch_state.is_some() => checked_probability(envelope.sample_p), + _ => Ok(1.0), + } } /// Envelope-wrapped state, or the bare state for producers that omit the @@ -61,7 +67,7 @@ pub fn carries_sketch_state(buffer: &[u8]) -> bool { /// A `SketchEnvelope{DdSketchState}` frame. Bare states are rejected. pub fn ddsketch_from_proto(buffer: &[u8]) -> Result { let (state, sample_p) = asap_sketch_codec::ddsketch_state(buffer)?; - unsampled("DDSketch", sample_p)?; + checked_probability(sample_p)?; if !(state.alpha > 0.0 && state.alpha < 1.0) { return Err(format!( "DDSketchState alpha {} out of range (expected 0 < alpha < 1)", @@ -177,7 +183,7 @@ fn read_uvarint(buf: &[u8]) -> Option<(u64, usize)> { let mut result: u64 = 0; for (i, &b) in buf.iter().enumerate() { let shift = 7 * i as u32; - if shift >= 64 { + if shift >= 64 || (shift == 63 && b & 0x7e != 0) { return None; } result |= u64::from(b & 0x7f) << shift; @@ -224,7 +230,7 @@ pub fn hll_from_proto(buffer: &[u8]) -> Result { sketch_envelope::SketchState::Hll(state) => Some(state), _ => None, })?; - unsampled("HLL", sample_p)?; + checked_probability(sample_p)?; if state.precision == 0 || state.precision > 20 { return Err(format!( "HyperLogLogState precision {} out of range (expected 1..=20)", @@ -269,8 +275,32 @@ pub fn hll_from_msgpack(buffer: &[u8]) -> Result { /// Apply an `HLLDelta` register frame (register-wise max) onto `sketch`. pub fn apply_hll_proto_delta(sketch: &mut HllSketch, buffer: &[u8]) -> Result<(), String> { + use asap_sketchlib::{proto::sketchlib::HllDelta, HllSketchDelta}; + let frame = HllDelta::decode(buffer).map_err(|e| format!("decode HLLDelta: {e}"))?; + let mut updates = Vec::new(); + let (mut previous, mut offset) = (0_u64, 0_usize); + while offset < frame.packed_updates.len() { + let (delta, used) = read_uvarint(&frame.packed_updates[offset..]) + .ok_or("HLLDelta: corrupt index varint")?; + offset += used; + let (value, used) = read_uvarint(&frame.packed_updates[offset..]) + .ok_or("HLLDelta: corrupt value varint")?; + offset += used; + let index = previous + .checked_add(delta) + .ok_or("HLLDelta: index overflow")?; + if index >= sketch.registers.len() as u64 || (!updates.is_empty() && delta == 0) { + return Err("HLLDelta: invalid register index".into()); + } + updates.push(( + u32::try_from(index).map_err(|_| "HLLDelta: index overflow")?, + u8::try_from(value).map_err(|_| "HLLDelta: register value exceeds u8")?, + )); + previous = index; + } + // Validate the whole frame before sketchlib mutates the cached registers. sketch - .apply_delta_bytes(buffer) + .apply_delta(&HllSketchDelta { updates }) .map_err(|e| format!("apply HLLDelta: {e}")) } @@ -340,7 +370,7 @@ pub fn cms_from_proto(buffer: &[u8]) -> Result { sketch_envelope::SketchState::CountMin(state) => Some(state), _ => None, })?; - unsampled("CountMin", sample_p)?; + checked_probability(sample_p)?; let (m, rows, cols) = matrix( "CountMinState", state.rows, @@ -397,7 +427,7 @@ pub fn cs_from_proto(buffer: &[u8]) -> Result { sketch_envelope::SketchState::CountSketch(state) => Some(state), _ => None, })?; - unsampled("CountSketch", sample_p)?; + checked_probability(sample_p)?; let (m, rows, cols) = matrix( "CountSketchState", state.rows, @@ -423,6 +453,14 @@ pub fn apply_cs_proto_delta(sketch: &mut CountSketch, buffer: &[u8]) -> Result<( &pb.cell_cols, &pb.d_counts, )?; + if pb.rows as usize != sketch.rows + || pb.cols as usize != sketch.cols + || cells + .iter() + .any(|(row, col, _)| *row as usize >= sketch.rows || *col as usize >= sketch.cols) + { + return Err("CountSketchDelta dimensions or cell indices differ from cached base".into()); + } let delta = CountSketchDelta { rows: pb.rows, cols: pb.cols, @@ -646,9 +684,9 @@ mod tests { ); } - // Malformed shapes, wrong families and sampled frames are rejected. + // Malformed shapes and wrong families fail; valid sampled frames decode. #[test] - fn invalid_or_sampled_frames_are_rejected() { + fn invalid_frames_fail_and_sampled_frames_decode() { assert!(cms_from_proto(&cms_state(2, 4, vec![1]).encode_to_vec()).is_err()); assert!(cms_from_proto(&cms_state(0, 4, vec![]).encode_to_vec()).is_err()); assert!(cms_from_proto(&cms_state(64, 1 << 20, vec![]).encode_to_vec()).is_err()); @@ -661,9 +699,8 @@ mod tests { sketch_envelope::SketchState::CountMin(cms_state(1, 2, vec![1, 1])), 0.5, ); - assert!(cms_from_proto(&sampled) - .unwrap_err() - .contains("edge-sampled")); + assert!(cms_from_proto(&sampled).is_ok()); + assert_eq!(sample_probability(&sampled).unwrap(), 0.5); let mut dd = DdSketch::new(0.01); dd.update(1.0); let mut dd_state = @@ -674,8 +711,7 @@ mod tests { sketch_envelope::SketchState::Ddsketch(dd_state.clone()), 0.25 )) - .unwrap_err() - .contains("edge-sampled")); + .is_ok()); dd_state.alpha = 0.0; assert!(ddsketch_from_proto(&envelope( sketch_envelope::SketchState::Ddsketch(dd_state), @@ -749,6 +785,37 @@ mod tests { assert_eq!(&sketch.registers[..2], &[2, 3]); } + // Rejected deltas must leave the base unchanged even after valid leading updates. + #[test] + fn invalid_delta_is_atomic() { + let mut hll = HllSketch::new(HllVariant::Regular, 4); + let before = hll.to_msgpack().unwrap(); + let delta = pb::HllDelta { + packed_updates: vec![0, 3, 16, 5], + } + .encode_to_vec(); + assert!(apply_hll_proto_delta(&mut hll, &delta).is_err()); + assert_eq!(hll.to_msgpack().unwrap(), before); + } + + // Rejected matrix deltas do not retain valid leading cell updates. + #[test] + fn invalid_cs_delta_is_atomic() { + let mut cs = CountSketch::new(3, 8); + let before = cs.to_msgpack().unwrap(); + let delta = pb::CountSketchDelta { + rows: 3, + cols: 8, + cell_rows: vec![0, 3], + cell_cols: vec![0, 0], + d_counts: vec![2, 4], + ..Default::default() + } + .encode_to_vec(); + assert!(apply_cs_proto_delta(&mut cs, &delta).is_err()); + assert_eq!(cs.to_msgpack().unwrap(), before); + } + // KLL frames reconstruct from their level layout and reject bad layouts. #[test] fn kll_frames_follow_their_level_layout() { diff --git a/crates/asap_summary_state/src/stored_state/delta_apply.rs b/crates/asap_summary_state/src/stored_state/delta_apply.rs index dc1df764d..4e7c33a86 100644 --- a/crates/asap_summary_state/src/stored_state/delta_apply.rs +++ b/crates/asap_summary_state/src/stored_state/delta_apply.rs @@ -11,7 +11,7 @@ use asap_sketchlib::{ use super::decoders as d; use super::{SketchEncoding, SketchSampleState}; -use crate::univmon::UnivMonAccumulator; +use asap_physical_operators::summary_kernels::univmon::UnivMonAccumulator; /// Which sketch family a candidate is, and the parameters needed to /// *bootstrap an empty state* — required by the per-window-reset (PWR) @@ -113,6 +113,29 @@ fn decode_full( encoding: SketchEncoding, ) -> Result { use SketchEncoding::{MsgpackFull, ProtoFull, WeightedFrequencyV1}; + if encoding == SketchEncoding::SampledKernelV2 { + let (sample_p, sketch): (f64, Vec) = + rmp_serde::from_slice(bytes).map_err(|e| e.to_string())?; + return Ok(match kind { + DeltaSketchKind::DDSketch { .. } => SummaryState::Dd( + k::DDSketchAccumulator::from_sketch(d::ddsketch_from_msgpack(&sketch)?, sample_p) + .map_err(|e| e.to_string())?, + ), + DeltaSketchKind::Hll { .. } => SummaryState::Hll( + k::HllSketchAccumulator::from_sketch(d::hll_from_msgpack(&sketch)?, sample_p) + .map_err(|e| e.to_string())?, + ), + DeltaSketchKind::Cms { .. } => SummaryState::Cms( + k::CountMinSketchAccumulator::from_sketch(d::cms_from_msgpack(&sketch)?, sample_p) + .map_err(|e| e.to_string())?, + ), + DeltaSketchKind::CountSketch { .. } => SummaryState::CountSketch( + k::CountSketchAccumulator::from_sketch(d::cs_from_msgpack(&sketch)?, sample_p) + .map_err(|e| e.to_string())?, + ), + _ => return Err("sampled kernel frame has an unsupported sketch family".into()), + }); + } Ok(match (kind, encoding) { ( DeltaSketchKind::UnivMon { @@ -123,7 +146,12 @@ fn decode_full( }, MsgpackFull, ) => { - let state = UnivMonAccumulator::from_bytes(bytes).map_err(|e| e.to_string())?; + let state = asap_sketchlib::UnivMon::deserialize_from_bytes(bytes) + .map_err(|e| e.to_string()) + .and_then(|sketch| { + UnivMonAccumulator::from_sketch(sketch).map_err(|e| e.to_string()) + }) + .map_err(|e| e.to_string())?; if state.dimensions() != ( *heap_size as usize, @@ -136,15 +164,39 @@ fn decode_full( } SummaryState::UnivMon(state) } - (DeltaSketchKind::DDSketch { .. }, ProtoFull) => dd(d::ddsketch_from_proto(bytes)?), + (DeltaSketchKind::DDSketch { .. }, ProtoFull) => SummaryState::Dd( + k::DDSketchAccumulator::from_sketch( + d::ddsketch_from_proto(bytes)?, + d::sample_probability(bytes)?, + ) + .map_err(|e| e.to_string())?, + ), (DeltaSketchKind::DDSketch { .. }, MsgpackFull) => dd(d::ddsketch_from_msgpack(bytes)?), - (DeltaSketchKind::Hll { .. }, ProtoFull) => hll(d::hll_from_proto(bytes)?), + (DeltaSketchKind::Hll { .. }, ProtoFull) => SummaryState::Hll( + k::HllSketchAccumulator::from_sketch( + d::hll_from_proto(bytes)?, + d::sample_probability(bytes)?, + ) + .map_err(|e| e.to_string())?, + ), (DeltaSketchKind::Hll { .. }, MsgpackFull) => hll(d::hll_from_msgpack(bytes)?), (DeltaSketchKind::Kll { .. }, ProtoFull) => kll(d::kll_from_proto(bytes)?), (DeltaSketchKind::Kll { .. }, MsgpackFull) => kll(d::kll_from_msgpack(bytes)?), - (DeltaSketchKind::Cms { .. }, ProtoFull) => cms(d::cms_from_proto(bytes)?), + (DeltaSketchKind::Cms { .. }, ProtoFull) => SummaryState::Cms( + k::CountMinSketchAccumulator::from_sketch( + d::cms_from_proto(bytes)?, + d::sample_probability(bytes)?, + ) + .map_err(|e| e.to_string())?, + ), (DeltaSketchKind::Cms { .. }, MsgpackFull) => cms(d::cms_from_msgpack(bytes)?), - (DeltaSketchKind::CountSketch { .. }, ProtoFull) => cs(d::cs_from_proto(bytes)?), + (DeltaSketchKind::CountSketch { .. }, ProtoFull) => SummaryState::CountSketch( + k::CountSketchAccumulator::from_sketch( + d::cs_from_proto(bytes)?, + d::sample_probability(bytes)?, + ) + .map_err(|e| e.to_string())?, + ), (DeltaSketchKind::CountSketch { .. }, MsgpackFull) => cs(d::cs_from_msgpack(bytes)?), // The legacy heap wire format is msgpack-only. (DeltaSketchKind::CmsWithHeap { .. }, ProtoFull | MsgpackFull) => { @@ -163,7 +215,7 @@ fn decode_full( _ => FrequencyAlgorithm::CountSketch, }; let state = super::codec::frequency_state(bytes).map_err(|e| e.to_string())?; - let kernel = super::codec::frequency_kernel(&state).map_err(|e| e.to_string())?; + let kernel = &state; // The catalog's heap size is not carried here; matrix shape is. let (width, depth, _) = kernel.shape(); if kernel.algorithm() != expected || (width, depth) != (*cols, *rows) { @@ -176,19 +228,23 @@ fn decode_full( } fn dd(inner: DdSketch) -> SummaryState { - SummaryState::Dd(k::DDSketchAccumulator { inner }) + SummaryState::Dd(k::DDSketchAccumulator::from_sketch(inner, 1.0).expect("unsampled sketch")) } fn hll(inner: HllSketch) -> SummaryState { - SummaryState::Hll(k::HllSketchAccumulator { inner }) + SummaryState::Hll(k::HllSketchAccumulator::from_sketch(inner, 1.0).expect("unsampled sketch")) } fn kll(inner: KllSketch) -> SummaryState { SummaryState::Kll(k::DatasketchesKLLAccumulator { inner }) } fn cms(inner: CountMinSketch) -> SummaryState { - SummaryState::Cms(k::CountMinSketchAccumulator { inner }) + SummaryState::Cms( + k::CountMinSketchAccumulator::from_sketch(inner, 1.0).expect("unsampled sketch"), + ) } fn cs(inner: CountSketch) -> SummaryState { - SummaryState::CountSketch(k::CountSketchAccumulator { inner }) + SummaryState::CountSketch( + k::CountSketchAccumulator::from_sketch(inner, 1.0).expect("unsampled sketch"), + ) } fn cms_heap_state(inner: CountMinSketchWithHeap) -> SummaryState { SummaryState::CmsWithHeap(k::CountMinSketchWithHeapAccumulator { inner }) @@ -289,32 +345,6 @@ impl SummaryState { } } - /// Rewrap a kernel of the same variant as `self`. - fn same_variant(&self, state: &dyn AggregateCore) -> Option { - let any = state.as_any(); - Some(match self { - Self::UnivMon(_) => Self::UnivMon(any.downcast_ref::()?.clone()), - Self::Dd(_) => Self::Dd(any.downcast_ref::()?.clone()), - Self::Hll(_) => Self::Hll(any.downcast_ref::()?.clone()), - Self::Kll(_) => Self::Kll(any.downcast_ref::()?.clone()), - Self::Cms(_) => Self::Cms(any.downcast_ref::()?.clone()), - Self::CountSketch(_) => { - Self::CountSketch(any.downcast_ref::()?.clone()) - } - Self::CmsWithHeap(_) => Self::CmsWithHeap( - any.downcast_ref::()? - .clone(), - ), - Self::CountSketchWithHeap(_) => Self::CountSketchWithHeap( - any.downcast_ref::()? - .clone(), - ), - Self::WeightedFrequency(_) => { - Self::WeightedFrequency(any.downcast_ref::()?.clone()) - } - }) - } - /// Apply a delta-encoded frame. DDSketch bucket deltas and HLL register /// deltas apply in place; every other frame decodes independently and /// merges in. @@ -342,6 +372,7 @@ impl SummaryState { let other = fragment(DeltaSketchKind::DDSketch { alpha: 0.0 }, ProtoFull)?; self.merge_same_family(&other) } else { + sketch.merge_sample_p(1.0).map_err(|e| e.to_string())?; d::apply_ddsketch_proto_delta(&mut sketch.inner, bytes) .map_err(|e| format!("apply DDSketch proto bucket-delta: {e}")) } @@ -350,7 +381,10 @@ impl SummaryState { let other = fragment(DeltaSketchKind::DDSketch { alpha: 0.0 }, MsgpackFull)?; self.merge_same_family(&other) } - (Self::Hll(sketch), ProtoDelta) => d::apply_hll_proto_delta(&mut sketch.inner, bytes), + (Self::Hll(sketch), ProtoDelta) => { + sketch.merge_sample_p(1.0).map_err(|e| e.to_string())?; + d::apply_hll_proto_delta(&mut sketch.inner, bytes) + } (Self::Hll(_), _) => { let other = fragment(DeltaSketchKind::Hll { precision: 0 }, MsgpackFull)?; self.merge_same_family(&other) @@ -417,29 +451,38 @@ impl SummaryState { } } - /// The bucket TOTAL — sum of row 0 of the underlying matrix, what a bare - /// (no item key) frequency readout reads. `0.0` for other families. - pub fn total(&self) -> f64 { - let matrix = match self { - Self::Cms(c) => c.inner.sketch(), - Self::CountSketch(c) => c.inner.sketch().clone(), - Self::CmsWithHeap(h) => h.inner.sketch_matrix(), - Self::CountSketchWithHeap(h) => h.inner.sketch_matrix(), - _ => return 0.0, - }; - matrix - .first() - .map(|row| row.iter().copied().sum::()) - .unwrap_or(0.0) + /// Total update mass where the Planner kernel supports an unkeyed count. + pub fn total(&self) -> Option { + use planner_types::{post_asap::SketchQuery, pre_asap::ColumnRef}; + match self { + Self::Cms(c) => Some(c.total()), + Self::CmsWithHeap(c) => Some(c.total()), + Self::Dd(c) => c + .estimate(&SketchQuery::PointCount { + key: ColumnRef::SampleValue, + value: None, + }) + .ok(), + // Signed CountSketch projections do not preserve total update mass. + _ => None, + } } /// Per-item point estimate; `None` for families without an item universe. pub fn estimate(&self, key: &str) -> Option { match self { - Self::Cms(c) => Some(c.inner.estimate(key)), - Self::CountSketch(c) => Some(c.inner.estimate(key)), - Self::CmsWithHeap(h) => Some(h.inner.estimate(key)), - Self::CountSketchWithHeap(h) => Some(h.inner.estimate(key)), + Self::Cms(c) => { + Some(c.query_key(&crate::KeyByLabelValues::new_with_labels(vec![key.into()]))) + } + Self::CountSketch(c) => { + Some(c.query_key(&crate::KeyByLabelValues::new_with_labels(vec![key.into()]))) + } + Self::CmsWithHeap(h) => { + Some(h.query_key(&crate::KeyByLabelValues::new_with_labels(vec![key.into()]))) + } + Self::CountSketchWithHeap(h) => { + Some(h.query_key(&crate::KeyByLabelValues::new_with_labels(vec![key.into()]))) + } _ => None, } } @@ -484,13 +527,33 @@ impl SummaryState { /// Merge `other` into `self` with the Planner kernel's merge; both must be /// the same family and shape. pub fn merge_same_family(&mut self, other: &SummaryState) -> Result<(), String> { - let merged = self + let mut merged = self .kernel() .merge_with(other.kernel()) .map_err(|e| format!("merge {} state: {e}", self.family_name()))?; - *self = self - .same_variant(merged.as_ref()) - .ok_or("Planner merge changed the state family")?; + // Move the merged value into the existing variant without cloning its sketch. + macro_rules! replace { + ($state:expr, $ty:ty) => { + std::mem::swap( + $state, + merged + .as_any_mut() + .downcast_mut::<$ty>() + .ok_or("Planner merge changed the state family")?, + ) + }; + } + match self { + Self::UnivMon(s) => replace!(s, UnivMonAccumulator), + Self::Dd(s) => replace!(s, k::DDSketchAccumulator), + Self::Hll(s) => replace!(s, k::HllSketchAccumulator), + Self::Kll(s) => replace!(s, k::DatasketchesKLLAccumulator), + Self::Cms(s) => replace!(s, k::CountMinSketchAccumulator), + Self::CountSketch(s) => replace!(s, k::CountSketchAccumulator), + Self::CmsWithHeap(s) => replace!(s, k::CountMinSketchWithHeapAccumulator), + Self::CountSketchWithHeap(s) => replace!(s, k::CountSketchWithHeapAccumulator), + Self::WeightedFrequency(s) => replace!(s, WeightedFrequency), + } Ok(()) } @@ -643,7 +706,8 @@ fn visit_window_summary_states( } SketchEncoding::ProtoFull | SketchEncoding::MsgpackFull - | SketchEncoding::WeightedFrequencyV1 => { + | SketchEncoding::WeightedFrequencyV1 + | SketchEncoding::SampledKernelV2 => { // A Full (re)sets this window's base. rolling = Some(decode_full(&kind, &state.bytes, state.encoding)?); } @@ -680,6 +744,62 @@ mod tests { //! got wrong — fails the build. use super::*; + // Local full-frame tags route sampled state through the versioned decoder. + #[test] + fn local_sampled_frames_keep_probability() { + use crate::StoredState; + let states: Vec<(DeltaSketchKind, Box)> = vec![ + ( + DeltaSketchKind::DDSketch { alpha: 0.01 }, + Box::new(k::DDSketchAccumulator::from_sketch(DdSketch::new(0.01), 0.25).unwrap()), + ), + ( + DeltaSketchKind::Hll { precision: 4 }, + Box::new( + k::HllSketchAccumulator::from_sketch( + HllSketch::new(HllVariant::Regular, 4), + 0.25, + ) + .unwrap(), + ), + ), + ( + DeltaSketchKind::Cms { rows: 2, cols: 16 }, + Box::new( + k::CountMinSketchAccumulator::from_sketch(CountMinSketch::new(2, 16), 0.25) + .unwrap(), + ), + ), + ( + DeltaSketchKind::CountSketch { rows: 3, cols: 16 }, + Box::new( + k::CountSketchAccumulator::from_sketch(CountSketch::new(3, 16), 0.25).unwrap(), + ), + ), + ]; + for (kind, state) in states { + let encoding = SketchEncoding::full_frame_for(state.as_ref()); + assert_eq!(encoding, SketchEncoding::SampledKernelV2); + assert!(encoding.is_full()); + let frame = SketchSampleState { + bytes: state.encode().unwrap(), + encoding, + }; + let restored = cumulative_summary_state(&[(1000, &frame)], kind) + .unwrap() + .unwrap(); + let probability = match restored { + SummaryState::Dd(s) => s.sample_p(), + SummaryState::Hll(s) => s.sample_p(), + SummaryState::Cms(s) => s.sample_p(), + SummaryState::CountSketch(s) => s.sample_p(), + _ => panic!("wrong restored family"), + }; + assert_eq!(probability, 0.25); + assert!(decode_full(&kind, &frame.bytes, SketchEncoding::MsgpackFull).is_err()); + } + } + // A stored Planner heap window decodes only for its catalog family and shape, merges // with another window, and ranks items under their series keys. #[test] @@ -693,9 +813,7 @@ mod tests { .unwrap(); state.update(&[Value::Utf8("plain".into())], 1.0).unwrap(); SketchSampleState { - bytes: crate::stored_state::codec::frequency_kernel(&state) - .unwrap() - .to_bytes(), + bytes: state.to_bytes(), encoding: SketchEncoding::WeightedFrequencyV1, } }; @@ -880,10 +998,9 @@ mod tests { sk } - // A sampled DDSketch envelope on the delta channel is rejected rather - // than read as an empty bucket delta. + // Sampled envelope deltas preserve their scaled count. #[test] - fn sampled_ddsketch_envelope_on_delta_channel_is_rejected() { + fn sampled_ddsketch_envelope_on_delta_channel_scales_count() { use asap_sketchlib::proto::sketchlib::{sketch_envelope, SketchEnvelope}; use prost::Message; let sketch = dd_over(0.01, &[1., 2.]); @@ -897,10 +1014,19 @@ mod tests { } .encode_to_vec(); let mut rolling = DeltaSketchKind::DDSketch { alpha: 0.01 }.bootstrap_empty(); - assert!(rolling + rolling .apply_delta_bytes(&sampled, SketchEncoding::ProtoDelta) - .unwrap_err() - .contains("edge-sampled")); + .unwrap(); + assert_eq!( + rolling + .kernel() + .estimate(&planner_types::post_asap::SketchQuery::PointCount { + key: planner_types::pre_asap::ColumnRef::SampleValue, + value: None, + }) + .unwrap(), + 4.0 + ); } /// A full re-snapshot replaces its pane's earlier frames; cumulative diff --git a/crates/asap_summary_state/src/stored_state/mod.rs b/crates/asap_summary_state/src/stored_state/mod.rs index 0c7400fde..160a408a8 100644 --- a/crates/asap_summary_state/src/stored_state/mod.rs +++ b/crates/asap_summary_state/src/stored_state/mod.rs @@ -23,6 +23,8 @@ pub enum SketchEncoding { /// A complete Planner weighted-frequency heap state (sketchlib /// `WeightedFrequency` bytes); never a legacy integer heap frame. WeightedFrequencyV1, + /// Planner sketch bytes paired with their sampling probability. + SampledKernelV2, } impl SketchEncoding { @@ -30,7 +32,7 @@ impl SketchEncoding { pub fn is_full(self) -> bool { matches!( self, - Self::ProtoFull | Self::MsgpackFull | Self::WeightedFrequencyV1 + Self::ProtoFull | Self::MsgpackFull | Self::WeightedFrequencyV1 | Self::SampledKernelV2 ) } @@ -39,6 +41,20 @@ impl SketchEncoding { use asap_physical_operators::summary_kernels::weighted_frequency::WeightedFrequency; if state.as_any().is::() { Self::WeightedFrequencyV1 + } else if state + .as_any() + .is::() + || state + .as_any() + .is::() + || state + .as_any() + .is::() + || state + .as_any() + .is::() + { + Self::SampledKernelV2 } else { Self::MsgpackFull } diff --git a/crates/asap_summary_state/src/stored_state/native.rs b/crates/asap_summary_state/src/stored_state/native.rs index 4f47c73ce..42ede9e7e 100644 --- a/crates/asap_summary_state/src/stored_state/native.rs +++ b/crates/asap_summary_state/src/stored_state/native.rs @@ -44,6 +44,9 @@ enum StateCodec { KllMsgpackV1, DdMsgpackV1, HllMsgpackV1, + DdSampledV2, + HllSampledV2, + ExactAccumulatorV2, } fn invalid(message: impl ToString) -> Error { Error::Invalid(message.to_string()) @@ -54,13 +57,13 @@ impl StateCodec { let codec = if any.is::() { Self::WeightedFrequencyV1 } else if any.is::() { - Self::ExactAccumulatorV1 + Self::ExactAccumulatorV2 } else if any.is::() { Self::KllMsgpackV1 } else if any.is::() { - Self::DdMsgpackV1 + Self::DdSampledV2 } else if any.is::() { - Self::HllMsgpackV1 + Self::HllSampledV2 } else { return Err(invalid("physical summary has no persisted native codec")); }; @@ -69,15 +72,23 @@ impl StateCodec { fn decode(&self, bytes: &[u8]) -> Result, Error> { let tag = match self { Self::WeightedFrequencyV1 => "WeightedFrequency", - Self::ExactAccumulatorV1 => codec::EXACT_V1, + Self::ExactAccumulatorV2 => codec::EXACT_V2, + Self::ExactAccumulatorV1 => { + return Err(invalid("native ExactAccumulatorV1 is retired; requires V2")) + } Self::SumAccumulatorV1 => { return Err(invalid( "native codec SumAccumulatorV1 is retired and no longer decoded", )) } Self::KllMsgpackV1 => "DatasketchesKLLAccumulator", - Self::DdMsgpackV1 => "DDSketchAccumulator", - Self::HllMsgpackV1 => "HllSketchAccumulator", + Self::DdSampledV2 => "DDSketchAccumulatorV2", + Self::DdMsgpackV1 | Self::HllMsgpackV1 => { + return Err(invalid( + "native sampled sketch codec V1 is retired; requires V2", + )) + } + Self::HllSampledV2 => "HllSketchAccumulatorV2", }; codec::decode(tag, bytes).map(Arc::from).map_err(invalid) } @@ -185,6 +196,10 @@ impl AggregateCore for NativeSummaryOutput { fn as_any(&self) -> &dyn std::any::Any { self } + fn as_any_mut(&mut self) -> &mut dyn std::any::Any { + self + } + fn merge_with(&self, _: &dyn AggregateCore) -> Result, KernelError> { Err("native output snapshots require an explicit physical merge operator".into()) } diff --git a/crates/asap_summary_state/src/stored_state/readout.rs b/crates/asap_summary_state/src/stored_state/readout.rs index 24796ed55..5f0a27787 100644 --- a/crates/asap_summary_state/src/stored_state/readout.rs +++ b/crates/asap_summary_state/src/stored_state/readout.rs @@ -51,7 +51,9 @@ pub fn sketch_query_value(rs: &SummaryState, query: &SketchQuery) -> Result Ok(rs.total()), + } => rs.total().ok_or(Error::Unsupported( + "sketch does not preserve total update mass", + )), SketchQuery::PointCount { key: ColumnRef::Named(_) | ColumnRef::Qualified { .. }, value: Some(v), @@ -228,3 +230,23 @@ mod counter_tests { assert!(exact_readout(states(), Statistic::Min, &None, &none).is_err()); } } + +#[cfg(test)] +mod tests { + use super::*; + // Signed CountSketch counters must not be misreported as population mass. + #[test] + fn signed_sketch_bare_count_is_unsupported() { + let state = SummaryState::CountSketch( + asap_physical_operators::summary_kernels::CountSketchAccumulator::new(3, 16), + ); + assert!(sketch_query_value( + &state, + &SketchQuery::PointCount { + key: ColumnRef::SampleValue, + value: None + } + ) + .is_err()); + } +} diff --git a/crates/asap_summary_state/src/univmon.rs b/crates/asap_summary_state/src/univmon.rs deleted file mode 100644 index 58546e961..000000000 --- a/crates/asap_summary_state/src/univmon.rs +++ /dev/null @@ -1,199 +0,0 @@ -//! SHIM: stored UnivMon state. -//! -//! Planner's `UnivMonAccumulator` keeps its sketchlib `UnivMon` private and has -//! no byte form, so stored UnivMon bytes cannot be decoded into it. Delete this -//! module once Planner exposes `UnivMonAccumulator::from_sketch(UnivMon)` and -//! `UnivMonAccumulator::sketch(&self) -> &UnivMon` (or a byte codec). -use asap_physical_operators::{AggregateCore, KernelError}; -use asap_sketchlib::{DataInput, UnivMon}; -use planner_types::{post_asap::SketchQuery, pre_asap::ColumnRef}; - -#[derive(Debug, Clone)] -pub struct UnivMonAccumulator { - inner: UnivMon, -} - -impl UnivMonAccumulator { - pub fn new( - heap_size: usize, - rows: usize, - cols: usize, - layers: usize, - ) -> Result { - if heap_size == 0 || cols == 0 || !(1..=20).contains(&rows) || !(1..=64).contains(&layers) { - return Err("invalid UnivMon dimensions".into()); - } - rows.checked_mul(cols) - .and_then(|n| n.checked_mul(layers)) - .ok_or("UnivMon dimensions overflow")?; - Ok(Self { - inner: UnivMon::init_univmon(heap_size, rows, cols, layers), - }) - } - - /// Each non-NaN sample is one occurrence. Signed zero has one identity. - pub fn insert_sample(&mut self, value: f64) -> Result<(), KernelError> { - if value.is_nan() { - return Ok(()); - } - self.inner - .bucket_size - .checked_add(1) - .ok_or("UnivMon count overflow")?; - let bits = if value == 0.0 { 0 } else { value.to_bits() }; - self.inner.insert(&DataInput::U64(bits), 1); - Ok(()) - } - - pub fn from_bytes(bytes: &[u8]) -> Result { - let inner = UnivMon::deserialize_from_bytes(bytes) - .map_err(|e| format!("invalid UnivMon state: {e}"))?; - if !inner.accepts_standard_updates() { - return Err( - "terminal-mode UnivMon state cannot enter the standard-update accumulator".into(), - ); - } - Ok(Self { inner }) - } - - pub fn to_bytes(&self) -> Result, KernelError> { - Ok(self.inner.serialize_to_bytes()?) - } - - pub fn dimensions(&self) -> (usize, usize, usize, usize) { - ( - self.inner.heap_size, - self.inner.sketch_row, - self.inner.sketch_col, - self.inner.layer_size, - ) - } - - /// Empty the sketch in place, keeping its shape. - pub fn clear(&mut self) { - self.inner.free(); - } - - pub fn merge_in_place(&mut self, other: &Self) -> Result<(), KernelError> { - if self.dimensions() != other.dimensions() { - return Err("incompatible UnivMon dimensions".into()); - } - self.inner - .bucket_size - .checked_add(other.inner.bucket_size) - .ok_or("UnivMon count overflow")?; - self.inner.merge(&other.inner); - Ok(()) - } -} - -impl AggregateCore for UnivMonAccumulator { - fn clone_boxed_core(&self) -> Box { - Box::new(self.clone()) - } - fn as_any(&self) -> &dyn std::any::Any { - self - } - fn merge_with(&self, other: &dyn AggregateCore) -> Result, KernelError> { - let other = other - .as_any() - .downcast_ref::() - .ok_or("expected UnivMon state")?; - let mut merged = self.clone(); - merged.merge_in_place(other)?; - Ok(Box::new(merged)) - } - fn estimate(&self, query: &SketchQuery) -> Result { - Ok(match query { - SketchQuery::PointCount { - key: ColumnRef::SampleValue, - value: None, - } => self.inner.calc_l1(), - SketchQuery::Cardinality => self.inner.calc_card(), - SketchQuery::FrequencyL2 => self.inner.calc_l2(), - SketchQuery::FrequencyEntropy => self.inner.calc_entropy(), - _ => return Err("unsupported UnivMon readout".into()), - }) - } - fn approx_memory_bytes(&self) -> usize { - std::mem::size_of::().saturating_add( - self.inner.layer_size.saturating_mul( - self.inner - .sketch_row - .saturating_mul(self.inner.sketch_col) - .saturating_mul(16) - .saturating_add(self.inner.heap_size.saturating_mul(256)), - ), - ) - } -} - -#[cfg(test)] -mod tests { - use super::*; - - fn read(state: &dyn AggregateCore, query: SketchQuery) -> f64 { - state.estimate(&query).unwrap() - } - fn count() -> SketchQuery { - SketchQuery::PointCount { - key: ColumnRef::SampleValue, - value: None, - } - } - - // Duplicate samples affect frequency but not cardinality, including signed zero. - #[test] - fn shared_readouts_survive_serialization() { - let mut state = UnivMonAccumulator::new(32, 5, 1024, 4).unwrap(); - for value in [0.0, -0.0, 2.0, 2.0, f64::NAN] { - state.insert_sample(value).unwrap(); - } - let restored = UnivMonAccumulator::from_bytes(&state.to_bytes().unwrap()).unwrap(); - for query in [ - count(), - SketchQuery::Cardinality, - SketchQuery::FrequencyL2, - SketchQuery::FrequencyEntropy, - ] { - assert_eq!(read(&state, query.clone()), read(&restored, query)); - } - assert_eq!(read(&restored, count()), 4.0); - assert!((read(&restored, SketchQuery::Cardinality) - 2.0).abs() < 0.01); - assert!((read(&restored, SketchQuery::FrequencyL2) - 8.0f64.sqrt()).abs() < 0.01); - assert!((read(&restored, SketchQuery::FrequencyEntropy) - 1.0).abs() < 0.01); - } - - // Terminal-mode serialization is valid sketchlib state but not this accumulator's update domain. - #[test] - fn terminal_state_is_rejected_before_ingestion_or_merge() { - let mut state = UnivMon::init_univmon(4, 3, 16, 2); - state.fast_insert(&DataInput::U64(1), 1); - let bytes = state.serialize_to_bytes().unwrap(); - assert!(UnivMonAccumulator::from_bytes(&bytes).is_err()); - state.free(); - assert!(UnivMonAccumulator::from_bytes(&state.serialize_to_bytes().unwrap()).is_ok()); - } - - // Pane merge preserves overlapping keys and clearing removes the previous window. - #[test] - fn merge_and_clear_preserve_frequency_semantics() { - let mut left = UnivMonAccumulator::new(32, 5, 1024, 4).unwrap(); - let mut right = left.clone(); - for value in [1.0, 2.0] { - left.insert_sample(value).unwrap(); - } - for value in [2.0, 3.0] { - right.insert_sample(value).unwrap(); - } - let merged = left.merge_with(&right).unwrap(); - assert_eq!(read(merged.as_ref(), count()), 4.0); - assert!((read(merged.as_ref(), SketchQuery::Cardinality) - 3.0).abs() < 0.01); - left.clear(); - assert_eq!(read(&left, count()), 0.0); - assert_eq!(read(&left, SketchQuery::FrequencyEntropy), 0.0); - assert!(left - .merge_with(&UnivMonAccumulator::new(16, 5, 1024, 4).unwrap()) - .is_err()); - } -} diff --git a/crates/asap_types/src/precompute_plan.rs b/crates/asap_types/src/precompute_plan.rs index 25d96a267..f9bc7fa6f 100644 --- a/crates/asap_types/src/precompute_plan.rs +++ b/crates/asap_types/src/precompute_plan.rs @@ -196,6 +196,8 @@ pub enum StateEncoding { ExactAccumulatorV1, /// Persisted backend state with explicit Planner family and population layout. PlannerExactAccumulatorV1, + /// Compensated Planner exact state; V1 payloads are retired. + PlannerExactAccumulatorV2, ExactCounterAccumulatorV2, } @@ -1266,11 +1268,11 @@ pub(crate) fn state_encodings(family: &SummaryFamilyType) -> Vec _, ) => vec![ StateEncoding::ExactCounterAccumulatorV2, - StateEncoding::PlannerExactAccumulatorV1, + StateEncoding::PlannerExactAccumulatorV2, ], SummaryFamilyType::ExactAggregate(..) => vec![ StateEncoding::ExactAccumulatorV1, - StateEncoding::PlannerExactAccumulatorV1, + StateEncoding::PlannerExactAccumulatorV2, ], SummaryFamilyType::Sketch(kind, _) if matches!( diff --git a/crates/asap_types/src/semantic_fragment.rs b/crates/asap_types/src/semantic_fragment.rs index 1e1f009ba..01455d684 100644 --- a/crates/asap_types/src/semantic_fragment.rs +++ b/crates/asap_types/src/semantic_fragment.rs @@ -483,7 +483,7 @@ mod tests { } // A transformed value cannot share state identity with its source column. #[test] - fn value_expression_is_semantic_and_nonfinite_constants_are_rejected() { + fn value_expression_and_nonfinite_constants_have_stable_semantics() { use planner_types::pre_asap::{ProjectItem, ScalarValue}; let original = fixture("latency"); let expected = SummarySemanticFragment::from_dag(&original, original.root).unwrap(); @@ -515,7 +515,12 @@ mod tests { unreachable!() }; cols[0].expr = QueryExpr::Literal(ScalarValue::Float64(f64::NAN)); - assert!(SummarySemanticFragment::from_dag(&transformed, transformed.root).is_err()); + let fragment = SummarySemanticFragment::from_dag(&transformed, transformed.root).unwrap(); + let bytes = canonical_bytes(&fragment).unwrap(); + let restored: SummarySemanticFragment = serde_json::from_slice(&bytes).unwrap(); + restored.validate().unwrap(); + assert_eq!(canonical_bytes(&restored).unwrap(), bytes); + assert_ne!(canonical_bytes(&logged).unwrap(), bytes); } // Changing a downstream consumer cannot change the persisted input definition. diff --git a/data_plane/examples/univmon_erp_artifact.rs b/data_plane/examples/univmon_erp_artifact.rs index a64610076..3c922cd95 100644 --- a/data_plane/examples/univmon_erp_artifact.rs +++ b/data_plane/examples/univmon_erp_artifact.rs @@ -1,7 +1,7 @@ //! Measure readout-specific ERP evidence from finite JSONL evaluation data. //! This offline tool retains samples; the production backend does not. use asap_physical_operators::summary_kernels::hll_sketch::HllSketchAccumulator; -use asap_summary_state::univmon::UnivMonAccumulator; +use asap_physical_operators::summary_kernels::univmon::UnivMonAccumulator; use data_plane::storage_engines::types::{AggregateCore, StoredState}; use serde_json::{json, Value}; use std::collections::{BTreeMap, HashMap}; @@ -109,7 +109,12 @@ fn main() -> Result<(), Box> { .map_err(|e| e.to_string())?; } left.merge_in_place(&right).map_err(|e| e.to_string())?; - bytes = bytes.max(left.to_bytes().map_err(|e| e.to_string())?.len()); + bytes = bytes.max( + left.sketch() + .serialize_to_bytes() + .map_err(|e| e.to_string())? + .len(), + ); for (i, stat) in [ planner_types::post_asap::SketchQuery::Cardinality, planner_types::post_asap::SketchQuery::FrequencyL2, diff --git a/data_plane/src/drivers/ingest/otel.rs b/data_plane/src/drivers/ingest/otel.rs index 72a58e84c..88940ec9b 100644 --- a/data_plane/src/drivers/ingest/otel.rs +++ b/data_plane/src/drivers/ingest/otel.rs @@ -2474,24 +2474,41 @@ fn decode_modified_otlp_sketch_bytes( // 3 — ENCODING_MSGPACK (full sketch-core msgpack state) // 4 — ENCODING_MSGPACK_DELTA (MSGPACK diff; not yet wired) Ok(match (encoding, algorithm) { - (ENCODING_PROTO, SketchAlgorithm::DDSketch) => Box::new(k::DDSketchAccumulator { - inner: d::ddsketch_from_proto(bytes)?, - }), + (ENCODING_PROTO, SketchAlgorithm::DDSketch) => Box::new( + k::DDSketchAccumulator::from_sketch( + d::ddsketch_from_proto(bytes)?, + d::sample_probability(bytes)?, + ) + .expect("validated sampling probability"), + ), (ENCODING_PROTO, SketchAlgorithm::Kll) => Box::new(k::DatasketchesKLLAccumulator { inner: d::kll_from_proto(bytes)?, }), - (ENCODING_PROTO, SketchAlgorithm::Cms) => Box::new(k::CountMinSketchAccumulator { - inner: d::cms_from_proto(bytes)?, - }), - (ENCODING_PROTO, SketchAlgorithm::CountSketch) => Box::new(k::CountSketchAccumulator { - inner: d::cs_from_proto(bytes)?, - }), - (ENCODING_PROTO, SketchAlgorithm::Hll) => Box::new(k::HllSketchAccumulator { - inner: d::hll_from_proto(bytes)?, - }), - (ENCODING_MSGPACK, SketchAlgorithm::Cms) => Box::new(k::CountMinSketchAccumulator { - inner: d::cms_from_msgpack(bytes)?, - }), + (ENCODING_PROTO, SketchAlgorithm::Cms) => Box::new( + k::CountMinSketchAccumulator::from_sketch( + d::cms_from_proto(bytes)?, + d::sample_probability(bytes)?, + ) + .expect("validated sampling probability"), + ), + (ENCODING_PROTO, SketchAlgorithm::CountSketch) => Box::new( + k::CountSketchAccumulator::from_sketch( + d::cs_from_proto(bytes)?, + d::sample_probability(bytes)?, + ) + .expect("validated sampling probability"), + ), + (ENCODING_PROTO, SketchAlgorithm::Hll) => Box::new( + k::HllSketchAccumulator::from_sketch( + d::hll_from_proto(bytes)?, + d::sample_probability(bytes)?, + ) + .expect("validated sampling probability"), + ), + (ENCODING_MSGPACK, SketchAlgorithm::Cms) => Box::new( + k::CountMinSketchAccumulator::from_sketch(d::cms_from_msgpack(bytes)?, 1.0) + .expect("validated sampling probability"), + ), (ENCODING_MSGPACK, SketchAlgorithm::CountSketch) => { // A heap-bearing CountSketch full frame is the // `{sketch,topk_heap,heap_size}` envelope, which the plain @@ -2502,20 +2519,23 @@ fn decode_modified_otlp_sketch_bytes( Ok(heap) if !heap.topk_heap_items().is_empty() => { Box::new(k::CountSketchWithHeapAccumulator { inner: heap }) } - _ => Box::new(k::CountSketchAccumulator { - inner: d::cs_from_msgpack(bytes)?, - }), + _ => Box::new( + k::CountSketchAccumulator::from_sketch(d::cs_from_msgpack(bytes)?, 1.0) + .expect("validated sampling probability"), + ), } } (ENCODING_MSGPACK, SketchAlgorithm::Kll) => Box::new(k::DatasketchesKLLAccumulator { inner: d::kll_from_msgpack(bytes)?, }), - (ENCODING_MSGPACK, SketchAlgorithm::DDSketch) => Box::new(k::DDSketchAccumulator { - inner: d::ddsketch_from_msgpack(bytes)?, - }), - (ENCODING_MSGPACK, SketchAlgorithm::Hll) => Box::new(k::HllSketchAccumulator { - inner: d::hll_from_msgpack(bytes)?, - }), + (ENCODING_MSGPACK, SketchAlgorithm::DDSketch) => Box::new( + k::DDSketchAccumulator::from_sketch(d::ddsketch_from_msgpack(bytes)?, 1.0) + .expect("validated sampling probability"), + ), + (ENCODING_MSGPACK, SketchAlgorithm::Hll) => Box::new( + k::HllSketchAccumulator::from_sketch(d::hll_from_msgpack(bytes)?, 1.0) + .expect("validated sampling probability"), + ), (ENCODING_PROTO, other) => { return Err( format!("modified-OTLP PROTO decoding is not implemented for {other:?}").into(), @@ -2647,22 +2667,38 @@ pub(crate) fn apply_modified_otlp_delta_bytes( (ENCODING_PROTO_DELTA, SketchAlgorithm::DDSketch) => edit( existing, "DDSketchAccumulator", - |s: &mut k::DDSketchAccumulator| d::apply_ddsketch_proto_delta(&mut s.inner, bytes), + |s: &mut k::DDSketchAccumulator| { + // Bare deltas inherit the series' probability from its full frame. + s.merge_sample_p(1.0).map_err(|e| e.to_string())?; + d::apply_ddsketch_proto_delta(&mut s.inner, bytes) + }, ), (ENCODING_PROTO_DELTA, SketchAlgorithm::Hll) => edit( existing, "HllSketchAccumulator", - |s: &mut k::HllSketchAccumulator| d::apply_hll_proto_delta(&mut s.inner, bytes), + |s: &mut k::HllSketchAccumulator| { + // Bare deltas inherit the series' probability from its full frame. + s.merge_sample_p(1.0).map_err(|e| e.to_string())?; + d::apply_hll_proto_delta(&mut s.inner, bytes) + }, ), (ENCODING_PROTO_DELTA, SketchAlgorithm::CountSketch) => edit( existing, "CountSketchAccumulator", - |s: &mut k::CountSketchAccumulator| d::apply_cs_proto_delta(&mut s.inner, bytes), + |s: &mut k::CountSketchAccumulator| { + // Bare deltas inherit the series' probability from its full frame. + s.merge_sample_p(1.0).map_err(|e| e.to_string())?; + d::apply_cs_proto_delta(&mut s.inner, bytes) + }, ), (ENCODING_PROTO_DELTA, SketchAlgorithm::Cms) => edit( existing, "CountMinSketchAccumulator", - |s: &mut k::CountMinSketchAccumulator| d::apply_cms_proto_delta(&mut s.inner, bytes), + |s: &mut k::CountMinSketchAccumulator| { + // Bare deltas inherit the series' probability from its full frame. + s.merge_sample_p(1.0).map_err(|e| e.to_string())?; + d::apply_cms_proto_delta(&mut s.inner, bytes) + }, ), (ENCODING_PROTO_DELTA, other) => Err(format!( "PROTO_DELTA for sketch kind {other:?} is not yet supported; \ @@ -2694,21 +2730,16 @@ pub(crate) fn apply_modified_otlp_delta_bytes( } /// Apply `update` to the cached base of kernel type `T`. Planner kernels -/// expose no mutable downcast, so the base is copied, updated and replaced. -fn edit( +/// expose mutable downcasting so updates do not copy the cached sketch. +fn edit( existing: &mut Box, expected: &str, update: impl FnOnce(&mut T) -> Result<(), String>, ) -> Result<(), Box> { - let mut state = existing - .as_any() - .downcast_ref::() - .ok_or_else(|| { - format!("apply_modified_otlp_delta_bytes: existing accumulator is not a {expected}") - })? - .clone(); - update(&mut state)?; - *existing = Box::new(state); + let state = existing.as_any_mut().downcast_mut::().ok_or_else(|| { + format!("apply_modified_otlp_delta_bytes: existing accumulator is not a {expected}") + })?; + update(state)?; Ok(()) } @@ -3262,15 +3293,53 @@ mod dispatcher_tests { use asap_sketchlib::DdSketch; use asap_sketchlib::HllVariant; + // Bare deltas inherit the full frame's probability and mutate the existing kernel. + #[test] + fn sampled_base_delta_preserves_probability_and_allocation() { + use asap_otel_proto::sketchlib::v1::{DdSketchBucketDelta, DdSketchDelta}; + use planner_types::{post_asap::SketchQuery, pre_asap::ColumnRef}; + let mut acc: Box = Box::new( + DDSketchAccumulator::from_sketch(DdSketch::from_raw(0.01, vec![1], 0), 0.25).unwrap(), + ); + let address = acc.as_any().downcast_ref::().unwrap() as *const _; + let delta = DdSketchDelta { + buckets: vec![DdSketchBucketDelta { + index: 0, + d_count: 1, + }], + } + .encode_to_vec(); + apply_modified_otlp_delta_bytes( + SketchAlgorithm::DDSketch, + ENCODING_PROTO_DELTA, + &mut acc, + &delta, + ) + .unwrap(); + assert_eq!( + address, + acc.as_any().downcast_ref::().unwrap() as *const _ + ); + assert_eq!( + acc.estimate(&SketchQuery::PointCount { + key: ColumnRef::SampleValue, + value: None + }) + .unwrap(), + 8.0 + ); + } + #[test] fn apply_modified_otlp_delta_bytes_ddsketch_round_trip() { use asap_otel_proto::sketchlib::v1::{DdSketchBucketDelta, DdSketchDelta as PbDelta}; use prost::Message; // Base sketch represents the last full snapshot the agent sent. - let mut acc: Box = Box::new(DDSketchAccumulator { - inner: DdSketch::from_raw(0.01, vec![1, 2, 3], 0), - }); + let mut acc: Box = Box::new( + DDSketchAccumulator::from_sketch(DdSketch::from_raw(0.01, vec![1, 2, 3], 0), 1.0) + .expect("validated sampling probability"), + ); // The wire delta now carries only bucket deltas (tags 2-7 // reserved post ProjectASAP/sketchlib-go#243 / asap_sketchlib#57). diff --git a/data_plane/src/precompute_engine/ingest_handler.rs b/data_plane/src/precompute_engine/ingest_handler.rs index feba9a10e..fc721afcf 100644 --- a/data_plane/src/precompute_engine/ingest_handler.rs +++ b/data_plane/src/precompute_engine/ingest_handler.rs @@ -377,9 +377,9 @@ mod tests { // first PROTO-encoded frame had landed. Use a distinct key // so we're not racing any earlier tests. let series_key = "__name__=latency_ms,inst=a"; - let base = DDSketchAccumulator { - inner: DdSketch::from_raw(0.01, vec![1, 2, 3], 0), - }; + let base = + DDSketchAccumulator::from_sketch(DdSketch::from_raw(0.01, vec![1, 2, 3], 0), 1.0) + .expect("validated sampling probability"); state.sketch_snapshots.insert( series_key.to_string(), SnapshotCacheEntry { diff --git a/data_plane/src/precompute_engine/raw_dag.rs b/data_plane/src/precompute_engine/raw_dag.rs index 15d9999e4..27f3301e9 100644 --- a/data_plane/src/precompute_engine/raw_dag.rs +++ b/data_plane/src/precompute_engine/raw_dag.rs @@ -277,24 +277,6 @@ impl RawDagProgram { .map_err(|e| e.to_string())?; return Ok(Box::new(state)); } - // Planner's UnivMon has no stored codec; the store keeps the - // backend UnivMon shim. - if let SketchParams::UnivMon { - heap_size, - sketch_rows, - sketch_cols, - layers, - } = kind.params() - { - return asap_summary_state::univmon::UnivMonAccumulator::new( - *heap_size as usize, - *sketch_rows as usize, - *sketch_cols as usize, - *layers as usize, - ) - .map(|state| Box::new(state) as Box) - .map_err(|e| e.to_string()); - } } Ok( create_planner_accumulator(&self.family, &self.input, &self.grouping)? diff --git a/data_plane/src/precompute_engine/worker.rs b/data_plane/src/precompute_engine/worker.rs index a367da198..ceae5517b 100644 --- a/data_plane/src/precompute_engine/worker.rs +++ b/data_plane/src/precompute_engine/worker.rs @@ -2885,7 +2885,7 @@ mod tests { // positive-only is the realistic shape. s.update(*v); } - DDSketchAccumulator { inner: s } + DDSketchAccumulator::from_sketch(s, 1.0).expect("validated sampling probability") } /// Pinning test: a single-group, single-window sketch ingest must @@ -2957,7 +2957,7 @@ mod tests { assert_eq!(output.end_timestamp, 90_000); assert_eq!( acc.type_name(), - "DDSketchAccumulator", + "DDSketchAccumulatorV2", "persisted accumulator must round-trip as DDSketchAccumulator (not silently demoted)" ); @@ -3228,7 +3228,7 @@ mod tests { assert_eq!(output.end_timestamp, 30_000); assert_eq!( acc.type_name(), - "DDSketchAccumulator", + "DDSketchAccumulatorV2", "wall-clock-fallback-emitted accumulator must round-trip as DDSketchAccumulator" ); let dd = acc @@ -3970,7 +3970,7 @@ mod tests { let (output, acc) = &captured[0]; assert_eq!(output.start_timestamp, 0); assert_eq!(output.end_timestamp, 30_000); - assert_eq!(acc.type_name(), "DDSketchAccumulator"); + assert_eq!(acc.type_name(), "DDSketchAccumulatorV2"); let dd = acc .as_any() .downcast_ref::() diff --git a/data_plane/src/query_engines/asap_query_engine/raw_source.rs b/data_plane/src/query_engines/asap_query_engine/raw_source.rs index 3baf830e2..43ae31bfe 100644 --- a/data_plane/src/query_engines/asap_query_engine/raw_source.rs +++ b/data_plane/src/query_engines/asap_query_engine/raw_source.rs @@ -866,6 +866,7 @@ mod tests { kind: BinaryOpKind::Arithmetic(ArithmeticOpKind::Add), vector_match: None, }, + [false, false], ) .unwrap(); let program = CompiledPhysicalDag::from_operators( 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 8db617a03..6b7cc4746 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 @@ -1223,7 +1223,7 @@ mod tests { #[test] fn bound_univmon_merges_panes_for_four_readouts() { use crate::storage_engines::sketch_db::index::SketchEncoding; - use asap_summary_state::univmon::UnivMonAccumulator; + use asap_physical_operators::summary_kernels::univmon::UnivMonAccumulator; use asap_types::query_plan::{MaterializationBinding, PhysicalGrouping}; let index = SketchStore::new(); let fp = asap_types::PolicyFingerprint(701); @@ -1252,7 +1252,7 @@ mod tests { BTreeMap::from([("job".into(), "a".into())]), (start, start + 1000), SketchSampleState { - bytes: state.to_bytes().unwrap(), + bytes: state.sketch().serialize_to_bytes().unwrap(), encoding: SketchEncoding::MsgpackFull, }, ); @@ -1312,18 +1312,20 @@ mod tests { use asap_physical_operators::summary_kernels::DDSketchAccumulator; let mut sketch = asap_sketchlib::DdSketch::new(0.01); assert!(sketch_query_value( - &SummaryState::Dd(DDSketchAccumulator { - inner: sketch.clone(), - }), + &SummaryState::Dd( + DDSketchAccumulator::from_sketch(sketch.clone(), 1.0) + .expect("validated sampling probability") + ), &SketchQuery::Quantile { q: 0.9 } ) .is_err()); sketch.update(20.0); for q in [0.0, 0.5, 0.9, 1.0] { let value = sketch_query_value( - &SummaryState::Dd(DDSketchAccumulator { - inner: sketch.clone(), - }), + &SummaryState::Dd( + DDSketchAccumulator::from_sketch(sketch.clone(), 1.0) + .expect("validated sampling probability"), + ), &SketchQuery::Quantile { q }, ) .unwrap(); @@ -1333,16 +1335,20 @@ mod tests { assert!(sketch.quantile(0.9).unwrap() < 21.0); for (q, expected) in [(0.0, 20.0), (0.5, 30.0), (0.9, 38.0), (1.0, 40.0)] { let value = sketch_query_value( - &SummaryState::Dd(DDSketchAccumulator { - inner: sketch.clone(), - }), + &SummaryState::Dd( + DDSketchAccumulator::from_sketch(sketch.clone(), 1.0) + .expect("validated sampling probability"), + ), &SketchQuery::Quantile { q }, ) .unwrap(); assert!((value - expected).abs() <= expected * 0.01); } assert!(sketch_query_value( - &SummaryState::Dd(DDSketchAccumulator { inner: sketch }), + &SummaryState::Dd( + DDSketchAccumulator::from_sketch(sketch, 1.0) + .expect("validated sampling probability") + ), &SketchQuery::Quantile { q: f64::NAN } ) .is_err()); diff --git a/data_plane/src/storage_engines/sketch_db/index/maintenance.rs b/data_plane/src/storage_engines/sketch_db/index/maintenance.rs index 9a3b3002c..02d9bb5fd 100644 --- a/data_plane/src/storage_engines/sketch_db/index/maintenance.rs +++ b/data_plane/src/storage_engines/sketch_db/index/maintenance.rs @@ -1385,18 +1385,34 @@ mod tests { let frames = &rows[0].samples[&60_000]; assert_eq!(frames.len(), 1); let sketch = match frames[0].encoding { - SketchEncoding::MsgpackFull => DDSketchAccumulator { - inner: asap_summary_state::stored_state::decoders::ddsketch_from_msgpack( + SketchEncoding::SampledKernelV2 => { + let state = asap_summary_state::stored_state::codec::decode( + "DDSketchAccumulatorV2", + &frames[0].bytes, + ) + .unwrap(); + state + .as_any() + .downcast_ref::() + .unwrap() + .clone() + } + SketchEncoding::MsgpackFull => DDSketchAccumulator::from_sketch( + asap_summary_state::stored_state::decoders::ddsketch_from_msgpack( &frames[0].bytes, ) .unwrap(), - }, - SketchEncoding::ProtoFull => DDSketchAccumulator { - inner: asap_summary_state::stored_state::decoders::ddsketch_from_proto( + 1.0, + ) + .expect("validated sampling probability"), + SketchEncoding::ProtoFull => DDSketchAccumulator::from_sketch( + asap_summary_state::stored_state::decoders::ddsketch_from_proto( &frames[0].bytes, ) .unwrap(), - }, + 1.0, + ) + .expect("validated sampling probability"), other => panic!("unexpected derived encoding: {other:?}"), }; assert_eq!(sketch.inner.total_count(), 2); diff --git a/data_plane/src/storage_engines/sketch_db/index/mod.rs b/data_plane/src/storage_engines/sketch_db/index/mod.rs index abccdb791..19037238a 100644 --- a/data_plane/src/storage_engines/sketch_db/index/mod.rs +++ b/data_plane/src/storage_engines/sketch_db/index/mod.rs @@ -63,6 +63,7 @@ fn encoding_to_tag(enc: SketchEncoding) -> u8 { SketchEncoding::MsgpackDelta => t::MSGPACK_DELTA, SketchEncoding::NativeBatchV1 => t::NATIVE_BATCH_V1, SketchEncoding::WeightedFrequencyV1 => t::WEIGHTED_FREQUENCY_V1, + SketchEncoding::SampledKernelV2 => t::SAMPLED_KERNEL_V2, } } @@ -78,6 +79,7 @@ fn tag_to_encoding(tag: u8) -> SketchEncoding { t::MSGPACK_DELTA => SketchEncoding::MsgpackDelta, t::NATIVE_BATCH_V1 => SketchEncoding::NativeBatchV1, t::WEIGHTED_FREQUENCY_V1 => SketchEncoding::WeightedFrequencyV1, + t::SAMPLED_KERNEL_V2 => SketchEncoding::SampledKernelV2, // t::PROTO_FULL and t::UNKNOWN (legacy) both → Full. _ => SketchEncoding::ProtoFull, } @@ -94,7 +96,7 @@ fn reconstruct_exact_agg( bytes: &[u8], ) -> Result>, String> { use asap_summary_state::stored_state::codec; - if type_name != codec::EXACT_V1 && !codec::is_retired_exact(type_name) { + if type_name != codec::EXACT_V2 && !codec::is_retired_exact(type_name) { return Ok(None); } codec::decode(type_name, bytes) diff --git a/data_plane/src/storage_engines/sketch_db/persistence/part.rs b/data_plane/src/storage_engines/sketch_db/persistence/part.rs index c426948ef..ea3f31428 100644 --- a/data_plane/src/storage_engines/sketch_db/persistence/part.rs +++ b/data_plane/src/storage_engines/sketch_db/persistence/part.rs @@ -108,6 +108,8 @@ pub mod encoding_tag { pub const MSGPACK_DELTA: u8 = 4; pub const NATIVE_BATCH_V1: u8 = 5; pub const WEIGHTED_FREQUENCY_V1: u8 = 6; + /// Planner sketch bytes with sampling probability. + pub const SAMPLED_KERNEL_V2: u8 = 7; } /// One entry inside a decoded part. The `start_ts`/`end_ts`/`label` diff --git a/data_plane/src/storage_engines/sketch_db/query/window_merger.rs b/data_plane/src/storage_engines/sketch_db/query/window_merger.rs index 7da26ae24..d51f35772 100644 --- a/data_plane/src/storage_engines/sketch_db/query/window_merger.rs +++ b/data_plane/src/storage_engines/sketch_db/query/window_merger.rs @@ -130,6 +130,10 @@ mod tests { self } + fn as_any_mut(&mut self) -> &mut dyn std::any::Any { + self + } + fn merge_with( &self, other: &dyn AggregateCore, diff --git a/data_plane/tests/edge_sketch_codec.rs b/data_plane/tests/edge_sketch_codec.rs index 306e8bec5..9b881ec6d 100644 --- a/data_plane/tests/edge_sketch_codec.rs +++ b/data_plane/tests/edge_sketch_codec.rs @@ -36,7 +36,8 @@ fn ddsketch_bare_state_is_rejected_and_envelope_supports_query_readout() { assert!(asap_sketch_codec::reconstruct_ddsketch(&bare).is_err()); let (decoded, _) = asap_sketch_codec::reconstruct_ddsketch(&envelope).unwrap(); let accumulator = - asap_physical_operators::summary_kernels::DDSketchAccumulator { inner: decoded }; + asap_physical_operators::summary_kernels::DDSketchAccumulator::from_sketch(decoded, 1.0) + .expect("validated sampling probability"); let median = accumulator .estimate(&planner_types::post_asap::SketchQuery::Quantile { q: 0.5 }) .unwrap(); diff --git a/data_plane/tests/support/univmon_erp_process.rs b/data_plane/tests/support/univmon_erp_process.rs index 9ca9ea769..411926f28 100644 --- a/data_plane/tests/support/univmon_erp_process.rs +++ b/data_plane/tests/support/univmon_erp_process.rs @@ -1,5 +1,5 @@ use super::*; -use asap_summary_state::univmon::UnivMonAccumulator; +use asap_physical_operators::summary_kernels::univmon::UnivMonAccumulator; use control_plane::physical::erp::ErpShapeObserver; use data_plane::storage_engines::types::AggregateCore; @@ -48,7 +48,7 @@ fn measured_artifact() -> Value { } let other = panes[1].clone(); panes[0].merge_in_place(&other).unwrap(); - bytes = bytes.max(panes[0].to_bytes().unwrap().len()); + bytes = bytes.max(panes[0].sketch().serialize_to_bytes().unwrap().len()); for (i, stat) in [ planner_types::post_asap::SketchQuery::Cardinality, planner_types::post_asap::SketchQuery::FrequencyL2,