Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
163 changes: 162 additions & 1 deletion crates/asap_summary_state/src/stored_state/delta_apply.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<Vec<_>>()
.join(";")
}

/// `name{k="v",...}` with labels sorted and values escaped, as series keys are.
fn series_key(labels: &std::collections::BTreeMap<String, String>) -> 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::<Vec<_>>();
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.
Expand All @@ -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 {
Expand Down Expand Up @@ -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)?
Expand Down Expand Up @@ -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,
}
}
Expand Down Expand Up @@ -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::<crate::summary_kernels::weighted_frequency::PhysicalWeightedFrequency>()
.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()
Expand All @@ -535,6 +639,7 @@ impl SummaryState {
SummaryState::CountSketch(_) => "CountSketch",
SummaryState::CmsWithHeap(_) => "CmsWithHeap",
SummaryState::CountSketchWithHeap(_) => "CountSketchWithHeap",
SummaryState::WeightedFrequency(_) => "WeightedFrequency",
}
}
}
Expand Down Expand Up @@ -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)?);
}
Expand Down Expand Up @@ -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<u8> {
Expand Down
25 changes: 25 additions & 0 deletions crates/asap_summary_state/src/stored_state/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<crate::summary_kernels::weighted_frequency::WeightedFrequency>()
{
Self::WeightedFrequencyV1
} else {
Self::MsgpackFull
}
}
}
5 changes: 5 additions & 0 deletions crates/asap_summary_state/src/stored_state/readout.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,11 @@ pub fn sketch_query_value(rs: &SummaryState, query: &SketchQuery) -> Result<f64,
.query_statistic(statistic, &None, &Default::default())
.map_err(|_| Error::Unsupported("UnivMon readout failed"));
}
if matches!(rs, SummaryState::WeightedFrequency(_)) {
return Err(Error::Unsupported(
"Planner weighted-frequency heaps support TopK readout only",
));
}
match query {
SketchQuery::FrequencyL2 | SketchQuery::FrequencyEntropy => Err(Error::Unsupported(
"frequency moment readout requires UnivMon",
Expand Down
Loading
Loading