From 8272a7abc925f563c8c9eeb05194a1627282f37c Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 10:06:18 +0000 Subject: [PATCH 1/2] chore: repin ASAPPlanner to integration rev 344369e Co-Authored-By: Claude Opus 5.5 --- Cargo.lock | 12 ++++++------ Cargo.toml | 10 +++++----- 2 files changed, 11 insertions(+), 11 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 6d9b9eeb4..025a00fdf 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=e7c64ab20a18592c462357c9fff954892b0d4a08#e7c64ab20a18592c462357c9fff954892b0d4a08" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=344369e25c50e65dbaaff95b4ea628c31999cb59#344369e25c50e65dbaaff95b4ea628c31999cb59" 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=e7c64ab20a18592c462357c9fff954892b0d4a08#e7c64ab20a18592c462357c9fff954892b0d4a08" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=344369e25c50e65dbaaff95b4ea628c31999cb59#344369e25c50e65dbaaff95b4ea628c31999cb59" 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=e7c64ab20a18592c462357c9fff954892b0d4a08#e7c64ab20a18592c462357c9fff954892b0d4a08" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=344369e25c50e65dbaaff95b4ea628c31999cb59#344369e25c50e65dbaaff95b4ea628c31999cb59" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -396,7 +396,7 @@ dependencies = [ [[package]] name = "asap-physical-operators" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=e7c64ab20a18592c462357c9fff954892b0d4a08#e7c64ab20a18592c462357c9fff954892b0d4a08" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=344369e25c50e65dbaaff95b4ea628c31999cb59#344369e25c50e65dbaaff95b4ea628c31999cb59" dependencies = [ "asap-types", "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib?rev=5f03ccbd798ed5fec62bdd839bcb331123cab369)", @@ -410,12 +410,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=e7c64ab20a18592c462357c9fff954892b0d4a08#e7c64ab20a18592c462357c9fff954892b0d4a08" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=344369e25c50e65dbaaff95b4ea628c31999cb59#344369e25c50e65dbaaff95b4ea628c31999cb59" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=e7c64ab20a18592c462357c9fff954892b0d4a08#e7c64ab20a18592c462357c9fff954892b0d4a08" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=344369e25c50e65dbaaff95b4ea628c31999cb59#344369e25c50e65dbaaff95b4ea628c31999cb59" dependencies = [ "serde", "serde_json", diff --git a/Cargo.toml b/Cargo.toml index ecc80e668..e31c43bf8 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 = "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"] } @@ -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" } From 859b7d9b78095e5111fb5ff18396f1c5c483e4b8 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 10:06:18 +0000 Subject: [PATCH 2/2] fix: restore grouped TopK heap item labels from their stored group A topk by heap partitions by its group labels, so Planner's weighted frequency item identity excludes them. Stored-state readout now merges the stored group labels back into each series identity before rendering it, and install no longer rejects identity items that exclude labels. Co-Authored-By: Claude Opus 5.5 --- .../src/stored_state/delta_apply.rs | 57 ++++++++-- .../src/stored_state/readout.rs | 8 +- data_plane/src/precompute_engine/raw_dag.rs | 16 --- .../asap_query_engine/post_asap_readout.rs | 2 +- .../asap_query_engine/summary_executor.rs | 34 ++++-- .../asapquery_compatibility_process_e2e.rs | 102 +++++++++++++----- 6 files changed, 157 insertions(+), 62 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 1888966b5..a81e2b4f2 100644 --- a/crates/asap_summary_state/src/stored_state/delta_apply.rs +++ b/crates/asap_summary_state/src/stored_state/delta_apply.rs @@ -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 { use asap_physical_operators::values::Value; items .iter() @@ -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(), @@ -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> { + /// no heap at all. `group` is the stored group the state belongs to. + pub fn topk_items( + &self, + group: &std::collections::BTreeMap, + ) -> Option> { match self { SummaryState::CmsWithHeap(h) => Some( h.topk_heap_items() @@ -566,7 +579,7 @@ impl SummaryState { else { return None; }; - Some((heap_item_key(&row), score)) + Some((heap_item_key(&row, group), score)) }) .collect(), ), @@ -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, @@ -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::>() + }; + 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 { diff --git a/crates/asap_summary_state/src/stored_state/readout.rs b/crates/asap_summary_state/src/stored_state/readout.rs index d1ec426fe..6d29a2d0f 100644 --- a/crates/asap_summary_state/src/stored_state/readout.rs +++ b/crates/asap_summary_state/src/stored_state/readout.rs @@ -82,8 +82,12 @@ pub fn sketch_query_value(rs: &SummaryState, query: &SketchQuery) -> Result Result, Error> { - let mut items = rs.topk_items().ok_or(Error::Unsupported( +pub fn topk_ranked( + rs: &SummaryState, + k: usize, + group: &std::collections::BTreeMap, +) -> Result, 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", ))?; diff --git a/data_plane/src/precompute_engine/raw_dag.rs b/data_plane/src/precompute_engine/raw_dag.rs index db9707981..e920bcfdc 100644 --- a/data_plane/src/precompute_engine/raw_dag.rs +++ b/data_plane/src/precompute_engine/raw_dag.rs @@ -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(), diff --git a/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs b/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs index d4fb9f414..2e8609baa 100644 --- a/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs +++ b/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs @@ -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 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 f23963d78..3560f017b 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 @@ -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, state: &GroupState, query: &SketchQuery, ) -> Result { - self.readout(state, query) + self.readout(group, state, query) } pub fn merge_bound_states( @@ -618,6 +620,7 @@ impl QueryExecutionContext<'_> { fn readout( &self, + group: &BTreeMap, state: &GroupState, query: &SketchQuery, ) -> Result { @@ -625,9 +628,9 @@ impl QueryExecutionContext<'_> { 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) } } } @@ -640,6 +643,7 @@ fn readout_cumulative( kind: DeltaSketchKind, query: &SketchQuery, t1_ms: i64, + group: &BTreeMap, ) -> Result { let mut merged: Option = None; let mut latest_window_end: Option = None; @@ -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 { @@ -697,6 +701,7 @@ fn readout_per_window( kind: DeltaSketchKind, query: &SketchQuery, t0_ms: i64, + group: &BTreeMap, ) -> Result { let mut by_window: BTreeMap = BTreeMap::new(); // Tracked from RAW window-ends, before the `w_end < t0_ms` carry-in @@ -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::, _>>()?; Ok(SummaryValue::TopK(points, coverage)) } else { @@ -763,8 +768,12 @@ fn sketch_query_value( }, ) } -fn topk_ranked(state: &SummaryState, k: usize) -> Result, SummaryExecutorError> { - asap_summary_state::stored_state::readout::topk_ranked(state, k).map_err( +fn topk_ranked( + state: &SummaryState, + k: usize, + group: &BTreeMap, +) -> Result, 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) }, @@ -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"); @@ -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") }; diff --git a/data_plane/tests/asapquery_compatibility_process_e2e.rs b/data_plane/tests/asapquery_compatibility_process_e2e.rs index 7bb99e832..8e231adfd 100644 --- a/data_plane/tests/asapquery_compatibility_process_e2e.rs +++ b/data_plane/tests/asapquery_compatibility_process_e2e.rs @@ -600,18 +600,50 @@ fn erp_collector_kll_export(plan: &Value, end_ms: u64, raw: &[f64], sequence: u6 // an installed QueryPlan, retaining all three ranked identities and values. #[tokio::test] async fn registered_temporal_topk_cms_heap() { - registered_temporal_topk(planner_types::post_asap::SketchAlgorithm::CmsWithHeap).await; + registered_temporal_topk( + planner_types::post_asap::SketchAlgorithm::CmsWithHeap, + false, + ) + .await; } #[tokio::test] async fn registered_temporal_topk_count_sketch_heap() { - registered_temporal_topk(planner_types::post_asap::SketchAlgorithm::CountSketchWithHeap).await; + registered_temporal_topk( + planner_types::post_asap::SketchAlgorithm::CountSketchWithHeap, + false, + ) + .await; } -async fn registered_temporal_topk(algorithm: planner_types::post_asap::SketchAlgorithm) { +// A grouped CMS heap ranks each job's series separately and names each ranked +// item by its whole series, including the job label the heap partitions by. +#[tokio::test] +async fn registered_temporal_topk_by_cms_heap() { + registered_temporal_topk(planner_types::post_asap::SketchAlgorithm::CmsWithHeap, true).await; +} + +// A grouped CountSketch heap ranks each job's series under its whole series. +#[tokio::test] +async fn registered_temporal_topk_by_count_sketch_heap() { + registered_temporal_topk( + planner_types::post_asap::SketchAlgorithm::CountSketchWithHeap, + true, + ) + .await; +} + +async fn registered_temporal_topk( + algorithm: planner_types::post_asap::SketchAlgorithm, + grouped: bool, +) { use control_plane::physical::compiler::{BackendLocalPlanningInput, DeploymentPlanCompiler}; use planner_types::post_asap::{CompositionOperator, SketchQuery, SummaryFamilyType}; - const QUERY: &str = "topk(3, count_over_time(top_endpoint_qps[5s]))"; + let query_text = if grouped { + "topk by (job) (2, count_over_time(top_endpoint_qps[5s]))" + } else { + "topk(3, count_over_time(top_endpoint_qps[5s]))" + }; struct Evidence; impl asap_aware_mapping::AccuracyEvidenceProvider for Evidence { fn propagation_stats( @@ -637,12 +669,12 @@ async fn registered_temporal_topk(algorithm: planner_types::post_asap::SketchAlg )) .unwrap(); let mut entry = fixture["query_workload"]["repeating_queries"][5].clone(); - entry["query"] = QUERY.into(); + entry["query"] = query_text.into(); entry["requirements"]["accuracy"] = serde_json::json!({"explicit": {"EpsilonDelta": {"epsilon": 0.05, "delta": 0.05}}}); fixture["query_workload"]["repeating_queries"] = serde_json::json!([entry]); fixture["implementation"]["topk_evidence"] = serde_json::json!({ - QUERY: { + query_text: { "selected_lower_bound": 95.0, "excluded_upper_bound": 80.0, "interval_failure_probability": 0.001, "observed_at_unix_ms": 9500, "source": "deterministic-count-ranking-fixture" @@ -652,7 +684,7 @@ async fn registered_temporal_topk(algorithm: planner_types::post_asap::SketchAlg let (mut request, environment) = snapshot.into_physical_compilation_request().unwrap(); let query = &mut request.queries[0]; let expr = control_plane::query_parser::parse_query_expr_canonical( - QUERY, + query_text, query.accuracy_target.clone(), ) .unwrap(); @@ -798,24 +830,28 @@ async fn registered_temporal_topk(algorithm: planner_types::post_asap::SketchAlg .as_millis() as i64; let base = now - now.rem_euclid(5000) - 20000; let counts = [ - ("alpha", 100), - ("beta", 50), - ("gamma", 200), - ("delta", 75), - ("epsilon", 10), - ("zeta", 150), + ("alpha", "a", 100), + ("beta", "a", 50), + ("gamma", "a", 200), + ("delta", "b", 75), + ("epsilon", "b", 10), + ("zeta", "b", 150), ]; let samples = WriteRequest { timeseries: counts .iter() - .map(|(item, count)| { + .map(|(item, job, count)| { // Non-unit values distinguish count updates from accidental weighted sums. let points = (0..2) .flat_map(|window| { (0..*count).map(move |i| (base + window * 5000 + 10 + i * 20, 17.0)) }) .collect::>(); - series_with_labels("top_endpoint_qps", &[("endpoint", item)], &points) + series_with_labels( + "top_endpoint_qps", + &[("endpoint", item), ("job", job)], + &points, + ) }) .collect(), }; @@ -823,7 +859,7 @@ async fn registered_temporal_topk(algorithm: planner_types::post_asap::SketchAlg let watermark = WriteRequest { timeseries: vec![series_with_labels( "top_endpoint_qps", - &[("endpoint", "gamma")], + &[("endpoint", "gamma"), ("job", "a")], &[(base + 10500, 17.0)], )], }; @@ -836,7 +872,7 @@ async fn registered_temporal_topk(algorithm: planner_types::post_asap::SketchAlg let response: Value = client .get(format!("{backend}/api/v1/query")) .query(&[ - ("query", QUERY.to_string()), + ("query", query_text.to_string()), ("time", timestamp.to_string()), ]) .send() @@ -853,20 +889,34 @@ async fn registered_temporal_topk(algorithm: planner_types::post_asap::SketchAlg }) .await .expect("registered TopK must become warm within 30s"); - let expected = [("gamma", 200.0), ("zeta", 150.0), ("alpha", 100.0)]; + let expected: &[(&str, &str, f64)] = if grouped { + &[ + ("gamma", "a", 200.0), + ("alpha", "a", 100.0), + ("zeta", "b", 150.0), + ("delta", "b", 75.0), + ] + } else { + &[ + ("gamma", "a", 200.0), + ("zeta", "b", 150.0), + ("alpha", "a", 100.0), + ] + }; let assert_ranks = |response: &Value, range: bool| { assert_eq!(response["status"], "success", "{response}"); assert!(is_warm(response), "{response}"); let rows = response["data"]["result"].as_array().unwrap(); - assert_eq!(rows.len(), 3, "{response}"); - for (item, count) in expected { + assert_eq!(rows.len(), expected.len(), "{response}"); + for &(item, job, count) in expected { + let series = format!("top_endpoint_qps{{endpoint=\"{item}\",job=\"{job}\"}}"); let row = rows .iter() - .find(|row| { - row["metric"]["item"].as_str() - == Some(format!("top_endpoint_qps{{endpoint=\"{item}\"}}").as_str()) - }) - .unwrap_or_else(|| panic!("missing {item}: {response}")); + .find(|row| row["metric"]["item"].as_str() == Some(series.as_str())) + .unwrap_or_else(|| panic!("missing {series}: {response}")); + if grouped { + assert_eq!(row["metric"]["job"], job, "{response}"); + } let points = if range { let points = row["values"].as_array().unwrap(); assert_eq!(points.len(), 2); @@ -888,7 +938,7 @@ async fn registered_temporal_topk(algorithm: planner_types::post_asap::SketchAlg let range: Value = client .get(format!("{backend}/api/v1/query_range")) .query(&[ - ("query", QUERY.to_string()), + ("query", query_text.to_string()), ("start", timestamp.to_string()), ("end", (timestamp + 5.0).to_string()), ("step", "5".into()),