From 4905bff932cbe46d0626f6b487a7980adc5500a0 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 08:11:33 +0000 Subject: [PATCH 1/3] feat(storage): store and read Planner weighted-frequency heap windows A raw heap output built by Planner is stored as a WeightedFrequencyV1 frame (sketchlib WeightedFrequency bytes, persisted with its own encoding tag). Stored heap readout decodes such frames for their catalog family, merges them, and ranks items under the legacy heap key (a canonical label identity renders as its series key). Point-count readouts over these frames are reported unsupported rather than guessed. Co-Authored-By: Claude Opus 5.5 --- .../src/stored_state/delta_apply.rs | 146 +++++++++++++++++- .../src/stored_state/mod.rs | 3 + .../src/stored_state/readout.rs | 5 + .../storage_engines/sketch_db/index/mod.rs | 11 +- .../sketch_db/persistence/part.rs | 1 + 5 files changed, 164 insertions(+), 2 deletions(-) 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 3f638503c..02c76f0bd 100644 --- a/crates/asap_summary_state/src/stored_state/delta_apply.rs +++ b/crates/asap_summary_state/src/stored_state/delta_apply.rs @@ -205,10 +205,74 @@ fn decode_full( ) => Ok(SummaryState::CountSketchWithHeap( decode_cs_with_heap_from_msgpack(bytes)?, )), + ( + DeltaSketchKind::CmsWithHeap { .. } | DeltaSketchKind::CountSketchWithHeap { .. }, + SketchEncoding::WeightedFrequencyV1, + ) => { + use crate::summary_kernels::weighted_frequency::FrequencyAlgorithm; + let kernel = asap_sketchlib::WeightedFrequency::from_bytes(bytes) + .map_err(|e| format!("deserialize weighted frequency: {e:?}"))?; + let expected = match kind { + DeltaSketchKind::CmsWithHeap { .. } => FrequencyAlgorithm::Cms, + _ => FrequencyAlgorithm::CountSketch, + }; + if kernel.algorithm() != expected { + return Err("weighted frequency algorithm differs from installed catalog".into()); + } + Ok(SummaryState::WeightedFrequency( + rmp_serde::from_slice(&rmp_serde::to_vec(&kernel).map_err(|e| e.to_string())?) + .map_err(|e| e.to_string())?, + )) + } (_, e) => Err(format!("decode_full called with non-Full encoding {e:?}")), } } +/// Render a ranked heap item as the legacy heap key: item parts joined by +/// `;`, with a canonical label-set identity shown as its series key. +fn heap_item_key(items: &[asap_physical_operators::values::Value]) -> String { + use asap_physical_operators::values::Value; + items + .iter() + .map(|item| match item { + Value::Utf8(text) => { + asap_physical_operators::physical_planner::promql_rows::decode_series_identity(text) + .map(|labels| series_key(&labels)) + .unwrap_or_else(|_| text.to_string()) + } + Value::Float64(value) => value.to_string(), + Value::Int64(value) => value.to_string(), + Value::Bool(value) => value.to_string(), + other => format!("{other:?}"), + }) + .collect::>() + .join(";") +} + +/// `name{k="v",...}` with labels sorted and values escaped, as series keys are. +fn series_key(labels: &std::collections::BTreeMap) -> String { + let name = labels + .get("__name__") + .map(String::as_str) + .unwrap_or_default(); + let pairs = labels + .iter() + .filter(|(k, _)| k.as_str() != "__name__") + .map(|(k, v)| { + let escaped = v + .replace('\\', "\\\\") + .replace('"', "\\\"") + .replace('\n', "\\n"); + format!("{k}=\"{escaped}\"") + }) + .collect::>(); + if pairs.is_empty() { + name.to_owned() + } else { + format!("{name}{{{}}}", pairs.join(",")) + } +} + /// The reconstructed state one candidate sid contributes — either /// folded across a window (or several) via delta application, or merged /// in from another sid's own reconstruction. @@ -223,6 +287,8 @@ pub enum SummaryState { /// these are two variants holding two different sketchlib types. CmsWithHeap(CountMinSketchWithHeap), CountSketchWithHeap(CountSketchWithHeap), + /// Planner's weighted-frequency heap, read through its ranked rows. + WeightedFrequency(crate::summary_kernels::weighted_frequency::PhysicalWeightedFrequency), } impl SummaryState { @@ -399,6 +465,9 @@ impl SummaryState { sk.merge(&other) .map_err(|e| format!("merge CmsWithHeap delta: {e}")) } + SummaryState::WeightedFrequency(_) => { + Err("weighted frequency frames are complete states, not deltas".into()) + } SummaryState::CountSketchWithHeap(sk) => { let other = if encoding == SketchEncoding::MsgpackDelta { decode_cs_with_heap_from_msgpack_delta(bytes)? @@ -478,6 +547,18 @@ impl SummaryState { .map(|item| (item.key, item.value)) .collect(), ), + SummaryState::WeightedFrequency(h) => Some( + h.rows(usize::MAX) + .into_iter() + .filter_map(|mut row| { + let asap_physical_operators::values::Value::Float64(score) = row.pop()? + else { + return None; + }; + Some((heap_item_key(&row), score)) + }) + .collect(), + ), _ => None, } } @@ -517,6 +598,18 @@ impl SummaryState { (SummaryState::CountSketchWithHeap(a), SummaryState::CountSketchWithHeap(b)) => a .merge(b) .map_err(|e| format!("merge CountSketchWithHeap: {e}")), + (SummaryState::WeightedFrequency(a), SummaryState::WeightedFrequency(b)) => { + use asap_physical_operators::AggregateCore; + let merged = a + .merge_with(b) + .map_err(|e| format!("merge weighted frequency: {e}"))?; + *a = merged + .as_any() + .downcast_ref::() + .ok_or("weighted frequency merge changed state type")? + .clone(); + Ok(()) + } (a, _) => Err(format!( "SummaryState family mismatch in merge_same_family (self is {})", a.family_name() @@ -535,6 +628,7 @@ impl SummaryState { SummaryState::CountSketch(_) => "CountSketch", SummaryState::CmsWithHeap(_) => "CmsWithHeap", SummaryState::CountSketchWithHeap(_) => "CountSketchWithHeap", + SummaryState::WeightedFrequency(_) => "WeightedFrequency", } } } @@ -670,7 +764,9 @@ fn visit_window_summary_states( SketchEncoding::NativeBatchV1 => { return Err("native physical outputs require the bound native batch decoder".into()) } - SketchEncoding::ProtoFull | SketchEncoding::MsgpackFull => { + SketchEncoding::ProtoFull + | SketchEncoding::MsgpackFull + | SketchEncoding::WeightedFrequencyV1 => { // A Full (re)sets this window's base. rolling = Some(decode_full(&kind, &state.bytes, state.encoding)?); } @@ -751,6 +847,54 @@ mod tests { //! notably the SPARSE-register HLL handling the deleted dead decoder //! got wrong — fails the build. use super::*; + + // A stored Planner heap window decodes for its catalog family only, merges + // with another window, and ranks items under their series keys. + #[test] + fn weighted_frequency_frames_rank_items_by_series_key() { + use crate::summary_kernels::weighted_frequency::{ + FrequencyAlgorithm, PhysicalWeightedFrequency, WeightedFrequency, + }; + use crate::SerializableToSink; + use asap_physical_operators::values::Value; + let frame = |weight| { + let mut state = + PhysicalWeightedFrequency::new(FrequencyAlgorithm::Cms, 64, 3, 8).unwrap(); + let identity = r#"{"__name__":"m","endpoint":"a\"b"}"#; + state + .update(&[Value::Utf8(identity.into())], weight) + .unwrap(); + state.update(&[Value::Utf8("plain".into())], 1.0).unwrap(); + SketchSampleState { + bytes: WeightedFrequency(state).serialize_to_bytes(), + encoding: SketchEncoding::WeightedFrequencyV1, + } + }; + let (first, second) = (frame(3.0), frame(4.0)); + let heap = DeltaSketchKind::CmsWithHeap { + rows: 3, + cols: 64, + heap_size: 8, + }; + let state = cumulative_summary_state(&[(1000, &first), (2000, &second)], heap) + .unwrap() + .unwrap(); + let mut items = state.topk_items().unwrap(); + items.sort_by(|a, b| a.0.cmp(&b.0)); + assert_eq!( + items, + vec![ + (r#"m{endpoint="a\"b"}"#.to_string(), 7.0), + ("plain".to_string(), 2.0) + ] + ); + let other = DeltaSketchKind::CountSketchWithHeap { + rows: 3, + cols: 64, + heap_size: 8, + }; + assert!(cumulative_summary_state(&[(1000, &first)], other).is_err()); + } use asap_sketchlib::HllVariant; fn encode_dd(sk: &DdSketch) -> Vec { diff --git a/crates/asap_summary_state/src/stored_state/mod.rs b/crates/asap_summary_state/src/stored_state/mod.rs index 89ab95786..7e2d229a4 100644 --- a/crates/asap_summary_state/src/stored_state/mod.rs +++ b/crates/asap_summary_state/src/stored_state/mod.rs @@ -19,4 +19,7 @@ pub enum SketchEncoding { MsgpackDelta, /// Versioned typed physical output; never a legacy sketch frame. NativeBatchV1, + /// A complete Planner weighted-frequency heap state (sketchlib + /// `WeightedFrequency` bytes); never a legacy integer heap frame. + WeightedFrequencyV1, } diff --git a/crates/asap_summary_state/src/stored_state/readout.rs b/crates/asap_summary_state/src/stored_state/readout.rs index df13f33c6..d1ec426fe 100644 --- a/crates/asap_summary_state/src/stored_state/readout.rs +++ b/crates/asap_summary_state/src/stored_state/readout.rs @@ -23,6 +23,11 @@ pub fn sketch_query_value(rs: &SummaryState, query: &SketchQuery) -> Result Err(Error::Unsupported( "frequency moment readout requires UnivMon", 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 25270f1c6..c1b06a469 100644 --- a/data_plane/src/storage_engines/sketch_db/index/mod.rs +++ b/data_plane/src/storage_engines/sketch_db/index/mod.rs @@ -61,6 +61,7 @@ fn encoding_to_tag(enc: SketchEncoding) -> u8 { SketchEncoding::MsgpackFull => t::MSGPACK_FULL, SketchEncoding::MsgpackDelta => t::MSGPACK_DELTA, SketchEncoding::NativeBatchV1 => t::NATIVE_BATCH_V1, + SketchEncoding::WeightedFrequencyV1 => t::WEIGHTED_FREQUENCY_V1, } } @@ -75,6 +76,7 @@ fn tag_to_encoding(tag: u8) -> SketchEncoding { t::MSGPACK_FULL => SketchEncoding::MsgpackFull, t::MSGPACK_DELTA => SketchEncoding::MsgpackDelta, t::NATIVE_BATCH_V1 => SketchEncoding::NativeBatchV1, + t::WEIGHTED_FREQUENCY_V1 => SketchEncoding::WeightedFrequencyV1, // t::PROTO_FULL and t::UNKNOWN (legacy) both → Full. _ => SketchEncoding::ProtoFull, } @@ -3409,7 +3411,14 @@ impl SketchStore { window, SketchSampleState { bytes: accumulator.serialize_to_bytes(), - encoding: SketchEncoding::MsgpackFull, + // Planner heap states are not legacy msgpack heap frames. + encoding: if accumulator.as_any().is::< + asap_summary_state::summary_kernels::weighted_frequency::WeightedFrequency, + >() { + SketchEncoding::WeightedFrequencyV1 + } else { + SketchEncoding::MsgpackFull + }, }, ), AggKind::ExactAgg { .. } => self.append_precompute_with_binding( 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 526ef99fb..c426948ef 100644 --- a/data_plane/src/storage_engines/sketch_db/persistence/part.rs +++ b/data_plane/src/storage_engines/sketch_db/persistence/part.rs @@ -107,6 +107,7 @@ pub mod encoding_tag { pub const MSGPACK_FULL: u8 = 3; pub const MSGPACK_DELTA: u8 = 4; pub const NATIVE_BATCH_V1: u8 = 5; + pub const WEIGHTED_FREQUENCY_V1: u8 = 6; } /// One entry inside a decoded part. The `start_ts`/`end_ts`/`label` From 8ef38ecaf387353c8dae286a08bc810cb1494b41 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 08:11:33 +0000 Subject: [PATCH 2/3] refactor(precompute): build raw heaps with the Planner graph CMS/CountSketch heaps now take the same Planner execution as every other raw output, so the backend's per-sample item and weight interpreter is deleted. Co-Authored-By: Claude Opus 5.5 --- data_plane/src/precompute_engine/raw_dag.rs | 126 ++------------------ 1 file changed, 7 insertions(+), 119 deletions(-) diff --git a/data_plane/src/precompute_engine/raw_dag.rs b/data_plane/src/precompute_engine/raw_dag.rs index 1f10204cc..ba1684d65 100644 --- a/data_plane/src/precompute_engine/raw_dag.rs +++ b/data_plane/src/precompute_engine/raw_dag.rs @@ -1,14 +1,14 @@ //! Bind raw ingestion to a selected Planner producer and its raw dependency edge. //! The backend supplies one typed sample batch per pane; the Planner-compiled //! precompute graph owns every update, grouping and item computation. -use crate::storage_engines::types::{AggregateCore, KeyByLabelValues}; +use crate::storage_engines::types::AggregateCore; use asap_physical_operators::{ operators::Operator, physical_planner::{precompute, CompiledPhysicalDag, Source as PhysicalSource}, runtime::{Limits, RunContext, Scope}, values::{Batch, Value}, }; -use asap_summary_state::factory::{create_planner_accumulator, AccumulatorUpdater}; +use asap_summary_state::factory::create_planner_accumulator; use asap_types::physical_plan_codec::PhysicalPlanCodec; use asap_types::{executable_plan::BackendNodeBinding, PrecomputeMaterialization}; use planner_types::post_asap::{ @@ -32,7 +32,7 @@ pub struct RawDagProgram { pub grouping: GroupingStrategy, pub reduction: planner_types::pre_asap::Reduction, /// Planner's encoded precompute graph from the raw sample boundary to this - /// output. Decoded graphs are not `Send`, so each execution decodes it. + /// output; decoded graphs are not `Send` (see `decoded`). program: std::sync::Arc<[u8]>, source: u64, } @@ -201,9 +201,6 @@ impl RawDagProgram { grouping: grouping.clone(), reduction: reduction.clone(), }; - if program.heap() { - program.validate_heap_update()?; - } if let Some(old) = &selected { if old.family != program.family || old.input != program.input @@ -223,16 +220,6 @@ impl RawDagProgram { selected.ok_or_else(|| "raw materialization has no selected post-ASAP DAG producer".into()) } - /// Heap readout decodes only the backend heap kernel, not Planner's weighted - /// frequency state, so heaps stay on that kernel until it does. - fn heap(&self) -> bool { - matches!(&self.family, SummaryFamilyType::Sketch(kind, _) if matches!( - kind.algorithm(), - planner_types::post_asap::SketchAlgorithm::CmsWithHeap - | planner_types::post_asap::SketchAlgorithm::CountSketchWithHeap - )) - } - /// Execute the Planner graph over one pane's samples as one typed batch. /// Returns `None` when the graph admits no population from them. pub fn build<'a>( @@ -242,13 +229,6 @@ impl RawDagProgram { max_bytes: usize, ) -> Result>, BuildError> { let samples = pane_batch(samples); - if self.heap() { - let mut updater = self.heap_updater()?; - for (series, time, value) in samples { - self.apply(&mut *updater, series, value, time)?; - } - return Ok(Some(updater.take_accumulator())); - } // The router assigns one population per group, so one pane yields // at most one state. match self.execute(samples, pane, max_bytes)?.as_slice() { @@ -270,13 +250,6 @@ impl RawDagProgram { max_bytes: usize, ) -> Result<(), BuildError> { let samples = pane_batch(samples); - if self.heap() { - let updater = self.heap_updater()?; - for (_, _, value) in samples { - self.validate_sample(updater.as_ref(), value)?; - } - return Ok(()); - } if samples.is_empty() { return Ok(()); } @@ -290,7 +263,10 @@ impl RawDagProgram { /// The family's empty state, for a pane known to have no samples. pub fn empty_state(&self) -> Result, String> { - Ok(self.heap_updater()?.take_accumulator()) + Ok( + create_planner_accumulator(&self.family, &self.input, &self.grouping)? + .take_accumulator(), + ) } fn execute<'a>( @@ -338,94 +314,6 @@ impl RawDagProgram { Ok::<_, asap_physical_operators::Error>(rows) })?) } - - /// Heaps still use the kernel interpreter, so install only updates it evaluates. - fn validate_heap_update(&self) -> Result<(), String> { - if !matches!( - &self.input.weight, - SummaryInputExpr::Column(ColumnRef::SampleValue) | SummaryInputExpr::Constant(_) - ) { - return Err("raw heap weight expression is unsupported".into()); - } - fn item(expr: &SummaryInputExpr) -> bool { - match expr { - SummaryInputExpr::Column(ColumnRef::Named(_) | ColumnRef::SampleValue) => true, - SummaryInputExpr::Tuple(items) => items.iter().all(item), - SummaryInputExpr::EntityIdentity( - planner_types::post_asap::EntityIdentity::PromqlLabelSet { excluding }, - ) => excluding.is_empty(), - _ => false, - } - } - if !self.input.item.as_ref().is_some_and(item) { - return Err("raw heap item expression is unsupported".into()); - } - self.heap_updater().map(|_| ()) - } - - fn heap_updater(&self) -> Result, String> { - create_planner_accumulator(&self.family, &self.input, &self.grouping) - } - - fn validate_sample(&self, updater: &dyn AccumulatorUpdater, value: f64) -> Result { - let weight = match &self.input.weight { - SummaryInputExpr::Constant(c) => *c, - // The worker retains one previous value per series across pane rotation. - SummaryInputExpr::Column(_) => value, - _ => return Err("unsupported raw weight expression".into()), - }; - let scalar_frequency = asap_types::accumulator_spec::is_unit_sample_frequency(&self.input) - && !updater.is_keyed(); - let weight = if scalar_frequency { value } else { weight }; - updater.validate_single_input(weight)?; - Ok(weight) - } - - fn apply( - &self, - updater: &mut dyn AccumulatorUpdater, - series: &str, - value: f64, - timestamp: i64, - ) -> Result<(), String> { - let weight = self.validate_sample(updater, value)?; - if updater.is_keyed() { - let labels = super::worker::parse_labels_from_series_key(series); - fn eval( - expr: &SummaryInputExpr, - series: &str, - value: f64, - labels: &HashMap<&str, &str>, - ) -> Result, String> { - Ok(match expr { - SummaryInputExpr::EntityIdentity(_) => vec![series.to_owned()], - SummaryInputExpr::Column(ColumnRef::SampleValue) => vec![value.to_string()], - SummaryInputExpr::Column(ColumnRef::Named(name)) => vec![labels - .get(name.as_str()) - .map(|s| super::worker::decode_label_value(s).into_owned()) - .ok_or_else(|| format!("missing DAG item column {name}"))?], - SummaryInputExpr::Tuple(items) => items - .iter() - .map(|i| eval(i, series, value, labels)) - .collect::, _>>()? - .into_iter() - .flatten() - .collect(), - _ => return Err("unsupported raw item expression".into()), - }) - } - let item = self - .input - .item - .as_ref() - .ok_or("keyed DAG kernel requires an explicit item")?; - let key = KeyByLabelValues::new_with_labels(eval(item, series, value, &labels)?); - updater.update_keyed(&key, weight, timestamp); - } else { - updater.update_single(weight, timestamp); - } - Ok(()) - } } /// The complete label set of a canonical series key, including `__name__`. From c98b0a2d78d28497ca9549f0053abab2eba55f68 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 08:32:46 +0000 Subject: [PATCH 3/3] fix(storage): tag every Planner heap frame and give heaps a Planner empty state Revision views and derived publications choose a stored frame's encoding from its state, so a Planner heap is never re-tagged as a legacy msgpack heap; carry-in checks treat it as a full frame. Heap frames decode only for their catalog's matrix shape, identity items render as series keys only when they name `__name__`, and a heap's empty state is a Planner heap, so revisions install for heap outputs. Heap items that omit identity labels are rejected at install until readout can name their series. Co-Authored-By: Claude Opus 5.5 --- .../src/stored_state/delta_apply.rs | 33 +++++-- .../src/stored_state/mod.rs | 22 +++++ data_plane/src/precompute_engine/raw_dag.rs | 86 ++++++++++++++++++- .../sketch_db/index/maintenance.rs | 4 +- .../storage_engines/sketch_db/index/mod.rs | 29 ++----- 5 files changed, 140 insertions(+), 34 deletions(-) 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 02c76f0bd..1888966b5 100644 --- a/crates/asap_summary_state/src/stored_state/delta_apply.rs +++ b/crates/asap_summary_state/src/stored_state/delta_apply.rs @@ -212,12 +212,19 @@ fn decode_full( use crate::summary_kernels::weighted_frequency::FrequencyAlgorithm; let kernel = asap_sketchlib::WeightedFrequency::from_bytes(bytes) .map_err(|e| format!("deserialize weighted frequency: {e:?}"))?; - let expected = match kind { - DeltaSketchKind::CmsWithHeap { .. } => FrequencyAlgorithm::Cms, - _ => FrequencyAlgorithm::CountSketch, + let (expected, rows, cols) = match kind { + DeltaSketchKind::CmsWithHeap { rows, cols, .. } => { + (FrequencyAlgorithm::Cms, *rows, *cols) + } + DeltaSketchKind::CountSketchWithHeap { rows, cols, .. } => { + (FrequencyAlgorithm::CountSketch, *rows, *cols) + } + _ => unreachable!("matched heap kinds"), }; - if kernel.algorithm() != expected { - return Err("weighted frequency algorithm differs from installed catalog".into()); + // 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) { + return Err("weighted frequency shape differs from installed catalog".into()); } Ok(SummaryState::WeightedFrequency( rmp_serde::from_slice(&rmp_serde::to_vec(&kernel).map_err(|e| e.to_string())?) @@ -229,7 +236,8 @@ fn decode_full( } /// Render a ranked heap item as the legacy heap key: item parts joined by -/// `;`, with a canonical label-set identity shown as its series key. +/// `;`, with a canonical series identity (it names `__name__`) shown as its +/// series key. fn heap_item_key(items: &[asap_physical_operators::values::Value]) -> String { use asap_physical_operators::values::Value; items @@ -237,9 +245,12 @@ fn heap_item_key(items: &[asap_physical_operators::values::Value]) -> String { .map(|item| match item { Value::Utf8(text) => { asap_physical_operators::physical_planner::promql_rows::decode_series_identity(text) + .ok() + .filter(|labels| labels.contains_key("__name__")) .map(|labels| series_key(&labels)) - .unwrap_or_else(|_| text.to_string()) + .unwrap_or_else(|| text.to_string()) } + Value::Null => String::new(), Value::Float64(value) => value.to_string(), Value::Int64(value) => value.to_string(), Value::Bool(value) => value.to_string(), @@ -848,7 +859,7 @@ mod tests { //! got wrong — fails the build. use super::*; - // A stored Planner heap window decodes for its catalog family only, merges + // 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] fn weighted_frequency_frames_rank_items_by_series_key() { @@ -894,6 +905,12 @@ mod tests { heap_size: 8, }; assert!(cumulative_summary_state(&[(1000, &first)], other).is_err()); + let narrower = DeltaSketchKind::CmsWithHeap { + rows: 3, + cols: 32, + heap_size: 8, + }; + assert!(cumulative_summary_state(&[(1000, &first)], narrower).is_err()); } use asap_sketchlib::HllVariant; diff --git a/crates/asap_summary_state/src/stored_state/mod.rs b/crates/asap_summary_state/src/stored_state/mod.rs index 7e2d229a4..e45b851e0 100644 --- a/crates/asap_summary_state/src/stored_state/mod.rs +++ b/crates/asap_summary_state/src/stored_state/mod.rs @@ -23,3 +23,25 @@ pub enum SketchEncoding { /// `WeightedFrequency` bytes); never a legacy integer heap frame. WeightedFrequencyV1, } + +impl SketchEncoding { + /// A complete window state that needs no earlier base. + pub fn is_full(self) -> bool { + matches!( + self, + Self::ProtoFull | Self::MsgpackFull | Self::WeightedFrequencyV1 + ) + } + + /// Encoding of a stored sketch state written as one complete window. + pub fn full_frame_for(state: &dyn crate::AggregateCore) -> Self { + if state + .as_any() + .is::() + { + Self::WeightedFrequencyV1 + } else { + Self::MsgpackFull + } + } +} diff --git a/data_plane/src/precompute_engine/raw_dag.rs b/data_plane/src/precompute_engine/raw_dag.rs index ba1684d65..db9707981 100644 --- a/data_plane/src/precompute_engine/raw_dag.rs +++ b/data_plane/src/precompute_engine/raw_dag.rs @@ -192,6 +192,22 @@ impl RawDagProgram { { return Err("raw precompute graph does not read the bound raw source".into()); } + // Stored heap readout renders identity items as whole series + // keys; one that omits labels would not name its series. + fn partial_identity(expr: &SummaryInputExpr) -> bool { + match expr { + SummaryInputExpr::EntityIdentity( + planner_types::post_asap::EntityIdentity::PromqlLabelSet { excluding }, + ) => !excluding.is_empty(), + SummaryInputExpr::Tuple(items) => items.iter().any(partial_identity), + _ => false, + } + } + if input.item.as_ref().is_some_and(partial_identity) { + return Err( + "raw heap items that exclude identity labels are not readable".into(), + ); + } let program = Self { source, program: compiled.encode().map_err(|e| e.to_string())?.into(), @@ -261,8 +277,38 @@ impl RawDagProgram { self.execute(samples, bounds, max_bytes).map(|_| ()) } - /// The family's empty state, for a pane known to have no samples. + /// The family's empty state, for a pane known to have no samples. Heaps + /// are Planner weighted-frequency states, as `build` produces. pub fn empty_state(&self) -> Result, String> { + use asap_summary_state::summary_kernels::weighted_frequency::{ + FrequencyAlgorithm, PhysicalWeightedFrequency, WeightedFrequency, + }; + use planner_types::post_asap::SketchParams; + if let SummaryFamilyType::Sketch(kind, _) = &self.family { + let heap = match kind.params() { + SketchParams::CmsWithHeap { + width, + depth, + heap_size, + } => Some((FrequencyAlgorithm::Cms, width, depth, heap_size)), + SketchParams::CountSketchWithHeap { + width, + depth, + heap_size, + } => Some((FrequencyAlgorithm::CountSketch, width, depth, heap_size)), + _ => None, + }; + if let Some((algorithm, width, depth, heap_size)) = heap { + let state = PhysicalWeightedFrequency::new( + algorithm, + *width as usize, + *depth as usize, + *heap_size as usize, + ) + .map_err(|e| e.to_string())?; + return Ok(Box::new(WeightedFrequency(state))); + } + } Ok( create_planner_accumulator(&self.family, &self.input, &self.grouping)? .take_accumulator(), @@ -371,3 +417,41 @@ fn decoded(encoded: &std::sync::Arc<[u8]>) -> Result encoding_to_tag(SketchEncoding::MsgpackFull), + AggKind::Sketch { .. } => { + encoding_to_tag(SketchEncoding::full_frame_for(state)) + } AggKind::ExactAgg { .. } => 0, } }, 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 c1b06a469..c116a379d 100644 --- a/data_plane/src/storage_engines/sketch_db/index/mod.rs +++ b/data_plane/src/storage_engines/sketch_db/index/mod.rs @@ -1915,10 +1915,7 @@ impl SketchStore { let Some(s) = payload.as_sketch() else { continue; }; - if !matches!( - s.encoding, - SketchEncoding::ProtoFull | SketchEncoding::MsgpackFull - ) { + if !s.encoding.is_full() { continue; } let w_end = win.1 as i64; @@ -2128,13 +2125,7 @@ impl SketchStore { }) .unwrap_or(false); let has_base_before = samples.iter().any(|(w_end, frames)| { - *w_end < start_unix_ms as i64 - && frames.iter().any(|s| { - matches!( - s.encoding, - SketchEncoding::ProtoFull | SketchEncoding::MsgpackFull - ) - }) + *w_end < start_unix_ms as i64 && frames.iter().any(|s| s.encoding.is_full()) }); earliest_is_delta && !has_base_before }) @@ -2161,10 +2152,7 @@ impl SketchStore { continue; }; let encoding = tag_to_encoding(entry.encoding_tag); - if !matches!( - encoding, - SketchEncoding::ProtoFull | SketchEncoding::MsgpackFull - ) { + if !encoding.is_full() { continue; } let label_map = Self::rebuild_label_map(&keys, &entry.label); @@ -3411,14 +3399,7 @@ impl SketchStore { window, SketchSampleState { bytes: accumulator.serialize_to_bytes(), - // Planner heap states are not legacy msgpack heap frames. - encoding: if accumulator.as_any().is::< - asap_summary_state::summary_kernels::weighted_frequency::WeightedFrequency, - >() { - SketchEncoding::WeightedFrequencyV1 - } else { - SketchEncoding::MsgpackFull - }, + encoding: SketchEncoding::full_frame_for(accumulator), }, ), AggKind::ExactAgg { .. } => self.append_precompute_with_binding( @@ -6936,7 +6917,7 @@ impl SketchStore { (record.start_ms, record.end_ms), SketchSampleState { bytes: state.serialize_to_bytes(), - encoding: SketchEncoding::MsgpackFull, + encoding: SketchEncoding::full_frame_for(state.as_ref()), }, ), AggKind::ExactAgg { .. } => view.append_precompute_with_binding(