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..1888966b5 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,85 @@ 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, rows, cols) = match kind { + DeltaSketchKind::CmsWithHeap { rows, cols, .. } => { + (FrequencyAlgorithm::Cms, *rows, *cols) + } + DeltaSketchKind::CountSketchWithHeap { rows, cols, .. } => { + (FrequencyAlgorithm::CountSketch, *rows, *cols) + } + _ => unreachable!("matched heap kinds"), + }; + // 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())?) + .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 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 + .iter() + .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()) + } + Value::Null => String::new(), + 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 +298,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 +476,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 +558,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 +609,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 +639,7 @@ impl SummaryState { SummaryState::CountSketch(_) => "CountSketch", SummaryState::CmsWithHeap(_) => "CmsWithHeap", SummaryState::CountSketchWithHeap(_) => "CountSketchWithHeap", + SummaryState::WeightedFrequency(_) => "WeightedFrequency", } } } @@ -670,7 +775,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 +858,60 @@ 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 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() { + 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()); + let narrower = DeltaSketchKind::CmsWithHeap { + rows: 3, + cols: 32, + heap_size: 8, + }; + assert!(cumulative_summary_state(&[(1000, &first)], narrower).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..e45b851e0 100644 --- a/crates/asap_summary_state/src/stored_state/mod.rs +++ b/crates/asap_summary_state/src/stored_state/mod.rs @@ -19,4 +19,29 @@ 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, +} + +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/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/precompute_engine/raw_dag.rs b/data_plane/src/precompute_engine/raw_dag.rs index 1f10204cc..db9707981 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, } @@ -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(), @@ -201,9 +217,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 +236,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 +245,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 +266,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(()); } @@ -288,9 +277,42 @@ 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> { - Ok(self.heap_updater()?.take_accumulator()) + 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(), + ) } fn execute<'a>( @@ -338,94 +360,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__`. @@ -483,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 25270f1c6..c116a379d 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, } @@ -1913,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; @@ -2126,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 }) @@ -2159,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); @@ -3409,7 +3399,7 @@ impl SketchStore { window, SketchSampleState { bytes: accumulator.serialize_to_bytes(), - encoding: SketchEncoding::MsgpackFull, + encoding: SketchEncoding::full_frame_for(accumulator), }, ), AggKind::ExactAgg { .. } => self.append_precompute_with_binding( @@ -6927,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( 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`