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
12 changes: 6 additions & 6 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

10 changes: 5 additions & 5 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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 = "e7c64ab20a18592c462357c9fff954892b0d4a08" }
asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "e7c64ab20a18592c462357c9fff954892b0d4a08" }
asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "e7c64ab20a18592c462357c9fff954892b0d4a08" }
asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "e7c64ab20a18592c462357c9fff954892b0d4a08" }
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" }

# Shared external deps (used by 2+ crates)
serde = { version = "1.0", features = ["derive"] }
Expand All @@ -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 = "e7c64ab20a18592c462357c9fff954892b0d4a08" }
asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "344369e25c50e65dbaaff95b4ea628c31999cb59" }
asap_sketch_codec = { path = "crates/asap_sketch_codec" }
asap_summary_state = { path = "crates/asap_summary_state" }
asap_types = { path = "crates/asap_types" }
Expand Down
57 changes: 50 additions & 7 deletions crates/asap_summary_state/src/stored_state/delta_apply.rs
Original file line number Diff line number Diff line change
Expand Up @@ -237,8 +237,12 @@ fn decode_full(

/// 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 {
/// series key. A grouped heap's identity omits the labels it is partitioned
/// by; those are stored as the heap's group labels and restored here.
fn heap_item_key(
items: &[asap_physical_operators::values::Value],
group: &std::collections::BTreeMap<String, String>,
) -> String {
use asap_physical_operators::values::Value;
items
.iter()
Expand All @@ -247,7 +251,13 @@ fn heap_item_key(items: &[asap_physical_operators::values::Value]) -> String {
asap_physical_operators::physical_planner::promql_rows::decode_series_identity(text)
.ok()
.filter(|labels| labels.contains_key("__name__"))
.map(|labels| series_key(&labels))
.map(|mut labels| {
// An empty group value is a series without that label.
for (name, value) in group.iter().filter(|(_, v)| !v.is_empty()) {
labels.entry(name.clone()).or_insert_with(|| value.clone());
}
series_key(&labels)
})
.unwrap_or_else(|| text.to_string())
}
Value::Null => String::new(),
Expand Down Expand Up @@ -543,8 +553,11 @@ impl SummaryState {
/// `None` for anything other than a heap-bearing state — the
/// heap-less Frequency states (`Cms`/`CountSketch`) carry no item
/// universe to enumerate, and the quantile/cardinality states have
/// no heap at all.
pub fn topk_items(&self) -> Option<Vec<(String, f64)>> {
/// no heap at all. `group` is the stored group the state belongs to.
pub fn topk_items(
&self,
group: &std::collections::BTreeMap<String, String>,
) -> Option<Vec<(String, f64)>> {
match self {
SummaryState::CmsWithHeap(h) => Some(
h.topk_heap_items()
Expand All @@ -566,7 +579,7 @@ impl SummaryState {
else {
return None;
};
Some((heap_item_key(&row), score))
Some((heap_item_key(&row, group), score))
})
.collect(),
),
Expand Down Expand Up @@ -890,7 +903,7 @@ mod tests {
let state = cumulative_summary_state(&[(1000, &first), (2000, &second)], heap)
.unwrap()
.unwrap();
let mut items = state.topk_items().unwrap();
let mut items = state.topk_items(&Default::default()).unwrap();
items.sort_by(|a, b| a.0.cmp(&b.0));
assert_eq!(
items,
Expand All @@ -912,6 +925,36 @@ mod tests {
};
assert!(cumulative_summary_state(&[(1000, &first)], narrower).is_err());
}

// A grouped heap's item identity omits its partition labels; readout
// restores non-empty ones from the stored group and keeps identity values.
#[test]
fn weighted_frequency_items_restore_group_labels() {
use crate::summary_kernels::weighted_frequency::{
FrequencyAlgorithm, PhysicalWeightedFrequency,
};
use asap_physical_operators::values::Value;
let mut state = PhysicalWeightedFrequency::new(FrequencyAlgorithm::Cms, 64, 3, 8).unwrap();
let identity = r#"{"__name__":"m","endpoint":"a"}"#;
state.update(&[Value::Utf8(identity.into())], 2.0).unwrap();
let state = SummaryState::WeightedFrequency(state);
let group = |pairs: &[(&str, &str)]| {
pairs
.iter()
.map(|(k, v)| (k.to_string(), v.to_string()))
.collect::<std::collections::BTreeMap<_, _>>()
};
assert_eq!(
state.topk_items(&group(&[("job", "j")])).unwrap(),
vec![(r#"m{endpoint="a",job="j"}"#.to_string(), 2.0)]
);
assert_eq!(
state
.topk_items(&group(&[("job", ""), ("endpoint", "other")]))
.unwrap(),
vec![(r#"m{endpoint="a"}"#.to_string(), 2.0)]
);
}
use asap_sketchlib::HllVariant;

fn encode_dd(sk: &DdSketch) -> Vec<u8> {
Expand Down
8 changes: 6 additions & 2 deletions crates/asap_summary_state/src/stored_state/readout.rs
Original file line number Diff line number Diff line change
Expand Up @@ -82,8 +82,12 @@ pub fn sketch_query_value(rs: &SummaryState, query: &SketchQuery) -> Result<f64,
/// heap-less family (`Dd`/`Hll`/`Kll`/`Cms`/`CountSketch` -- no item
/// universe to rank), not for an empty heap (a heap-bearing family that
/// simply never received any updates yields `Ok(vec![])`, not an error).
pub fn topk_ranked(rs: &SummaryState, k: usize) -> Result<Vec<(String, f64)>, Error> {
let mut items = rs.topk_items().ok_or(Error::Unsupported(
pub fn topk_ranked(
rs: &SummaryState,
k: usize,
group: &std::collections::BTreeMap<String, String>,
) -> Result<Vec<(String, f64)>, Error> {
let mut items = rs.topk_items(group).ok_or(Error::Unsupported(
"TopK requires a heap-bearing family (CmsWithHeap/CountSketchWithHeap) -- \
this state's family carries no item universe to rank",
))?;
Expand Down
16 changes: 0 additions & 16 deletions data_plane/src/precompute_engine/raw_dag.rs
Original file line number Diff line number Diff line change
Expand Up @@ -192,22 +192,6 @@ 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(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -212,7 +212,7 @@ impl PhysicalQueryRuntime<'_> {
for (key, state) in groups {
let value = self
.context
.readout_bound(state, &query)
.readout_bound(key, state, &query)
.map_err(PhysicalNodeError::Store)?;
let (mut rows, row_coverage) = expand_item_readout(key, value, item_labels)?;
if self.language == control_plane::query_plan::QueryLanguage::MetricsQl
Expand Down
34 changes: 24 additions & 10 deletions data_plane/src/query_engines/asap_query_engine/summary_executor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -548,12 +548,14 @@ impl QueryExecutionContext<'_> {
Ok(result)
}

/// `group` is the state's stored group labels, as `read_bound_materialization` returns them.
pub fn readout_bound(
&self,
group: &BTreeMap<String, String>,
state: &GroupState,
query: &SketchQuery,
) -> Result<SummaryValue, SummaryExecutorError> {
self.readout(state, query)
self.readout(group, state, query)
}

pub fn merge_bound_states(
Expand Down Expand Up @@ -618,16 +620,17 @@ impl QueryExecutionContext<'_> {

fn readout(
&self,
group: &BTreeMap<String, String>,
state: &GroupState,
query: &SketchQuery,
) -> Result<SummaryValue, SummaryExecutorError> {
let GroupState::Sketch { entries, kind } = state else {
return Err(SummaryExecutorError::UnsupportedFamily);
};
if self.is_cumulative {
readout_cumulative(entries, *kind, query, self.t1_ms as i64)
readout_cumulative(entries, *kind, query, self.t1_ms as i64, group)
} else {
readout_per_window(entries, *kind, query, self.t0_ms as i64)
readout_per_window(entries, *kind, query, self.t0_ms as i64, group)
}
}
}
Expand All @@ -640,6 +643,7 @@ fn readout_cumulative(
kind: DeltaSketchKind,
query: &SketchQuery,
t1_ms: i64,
group: &BTreeMap<String, String>,
) -> Result<SummaryValue, SummaryExecutorError> {
let mut merged: Option<SummaryState> = None;
let mut latest_window_end: Option<i64> = None;
Expand Down Expand Up @@ -675,7 +679,7 @@ fn readout_cumulative(
let w_end = latest_window_end.unwrap_or(t1_ms);
if let SketchQuery::TopK { k } = query {
Ok(SummaryValue::TopK(
vec![(w_end, topk_ranked(&merged, *k)?)],
vec![(w_end, topk_ranked(&merged, *k, group)?)],
coverage,
))
} else {
Expand All @@ -697,6 +701,7 @@ fn readout_per_window(
kind: DeltaSketchKind,
query: &SketchQuery,
t0_ms: i64,
group: &BTreeMap<String, String>,
) -> Result<SummaryValue, SummaryExecutorError> {
let mut by_window: BTreeMap<i64, SummaryState> = BTreeMap::new();
// Tracked from RAW window-ends, before the `w_end < t0_ms` carry-in
Expand Down Expand Up @@ -738,7 +743,7 @@ fn readout_per_window(
if let SketchQuery::TopK { k } = query {
let points = by_window
.into_iter()
.map(|(w_end, rs)| topk_ranked(&rs, *k).map(|items| (w_end, items)))
.map(|(w_end, rs)| topk_ranked(&rs, *k, group).map(|items| (w_end, items)))
.collect::<Result<Vec<_>, _>>()?;
Ok(SummaryValue::TopK(points, coverage))
} else {
Expand All @@ -763,8 +768,12 @@ fn sketch_query_value(
},
)
}
fn topk_ranked(state: &SummaryState, k: usize) -> Result<Vec<(String, f64)>, SummaryExecutorError> {
asap_summary_state::stored_state::readout::topk_ranked(state, k).map_err(
fn topk_ranked(
state: &SummaryState,
k: usize,
group: &BTreeMap<String, String>,
) -> Result<Vec<(String, f64)>, SummaryExecutorError> {
asap_summary_state::stored_state::readout::topk_ranked(state, k, group).map_err(
|asap_summary_state::stored_state::readout::Error::Unsupported(reason)| {
SummaryExecutorError::Unsupported(reason)
},
Expand Down Expand Up @@ -1187,7 +1196,11 @@ mod tests {
};
let states = context.read_bound_materialization(&binding).unwrap();
let SummaryValue::Points(points, coverage) = context
.readout_bound(&states[0].1, &SketchQuery::Quantile { q: 0.5 })
.readout_bound(
&states[0].0,
&states[0].1,
&SketchQuery::Quantile { q: 0.5 },
)
.unwrap()
else {
panic!("expected points");
Expand Down Expand Up @@ -1267,8 +1280,9 @@ mod tests {
(SketchQuery::FrequencyL2, 6.0f64.sqrt()),
(SketchQuery::FrequencyEntropy, 1.5),
] {
let SummaryValue::Points(points, _) =
context.readout_bound(&states[0].1, &query).unwrap()
let SummaryValue::Points(points, _) = context
.readout_bound(&states[0].0, &states[0].1, &query)
.unwrap()
else {
panic!("expected scalar points")
};
Expand Down
Loading
Loading