From 12d774072d70ff8f5d5efc7c399d8659c441d5bd Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 15 May 2026 23:15:52 -0600 Subject: [PATCH] feat(query): InstantVectorElement carries per-element label_keys_override MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Mirrors `RangeVectorElement` — adds an optional `label_keys_override: Option>` field so per-element keys (synthesized by ASAP-tier `topk(...)` as `"item": `) surface through the Prometheus adapter on instant-vector responses, not just range-vector. Three-line change spanning three files: - `InstantVectorElement` gains the field + `with_label_keys_override` helper (query_result.rs). - `asap_tier_result_to_query_result`'s instant-vector branch stops discarding the BTreeMap keys (engine.rs:3443) — the same fix PR #256 applied to the range-vector branch. - `convert_query_result_to_prometheus` picks per-element override over the query-scoped `KeyByLabelNames`, mirroring `convert_range_result_to_prometheus` (http.rs). Tests 8+9 in `e2e_controller_plans_and_backend_serves.rs` tighten their `topk(3, top_endpoint_qps)` assertions to verify `metric.item == "gamma"` is present in at least one returned series — the strict invariant that previously couldn't be checked because the adapter dropped the synthesized key. Co-Authored-By: Claude Opus 4.7 (1M context) --- .../query_engines/asap_query_engine/engine.rs | 23 ++++++++------- data_plane/src/query_engines/query_result.rs | 22 +++++++++++++- data_plane/src/utils/http.rs | 19 +++++++++--- ...e2e_controller_plans_and_backend_serves.rs | 29 +++++++++++++------ 4 files changed, 68 insertions(+), 25 deletions(-) diff --git a/data_plane/src/query_engines/asap_query_engine/engine.rs b/data_plane/src/query_engines/asap_query_engine/engine.rs index ef088013..55c8db75 100644 --- a/data_plane/src/query_engines/asap_query_engine/engine.rs +++ b/data_plane/src/query_engines/asap_query_engine/engine.rs @@ -3441,21 +3441,22 @@ fn asap_tier_result_to_query_result( if !is_range_query { let mut elements: Vec = Vec::with_capacity(result.series.len()); for (label_values, samples) in result.series { - let (_keys, values): (Vec, Vec) = label_values.into_iter().unzip(); + // Mirror the range-vector branch: BTreeMap iteration is + // key-sorted, so `unzip` produces aligned (keys, values). + // Stash the keys in the per-element `label_keys_override` + // so the Prometheus adapter renders synthesized keys + // (notably ASAP-tier `topk`'s `"item"` key) instead of + // the empty `metric: {}` it would produce when the + // query-scoped `KeyByLabelNames` is empty. + let (keys, values): (Vec, Vec) = label_values.into_iter().unzip(); let labels = KeyByLabelValues::new_with_labels(values); // Take the latest sample (the reducer returns one per // window_end; for instant readout we want the most recent). - // `InstantVectorElement` doesn't carry a per-element - // `label_keys_override` today (only `RangeVectorElement` - // does, for the topk-`item`-key case) — labels render - // with whatever query-scoped `KeyByLabelNames` the - // serializer holds. That's correct for the cardinality - // shape that's the only instant-vector consumer at the - // moment; if a future instant-vector readout needs - // per-element key remapping, add the override field on - // `InstantVectorElement` then plumb `keys` here. if let Some((_, value)) = samples.into_iter().last() { - elements.push(InstantVectorElement::new(labels, value)); + elements.push( + InstantVectorElement::new(labels, value) + .with_label_keys_override(keys), + ); } } return QueryResult::vector(elements, now_ms); diff --git a/data_plane/src/query_engines/query_result.rs b/data_plane/src/query_engines/query_result.rs index 6eaafa48..ff47e26c 100644 --- a/data_plane/src/query_engines/query_result.rs +++ b/data_plane/src/query_engines/query_result.rs @@ -153,11 +153,31 @@ pub struct InstantVector { pub struct InstantVectorElement { pub labels: KeyByLabelValues, pub value: f64, + /// Optional per-element label-key override. Mirrors the field on + /// `RangeVectorElement` — when `Some`, the HTTP serializer uses + /// these keys for the PromQL response's `"metric"` object instead + /// of the query-scoped `KeyByLabelNames` argument. Used by ASAP-tier + /// `topk(...)` (whose reducer synthesizes an `"item"` key not + /// present in the query's group-by clause). `None` for everyone + /// else — the existing serializer path is unaffected. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub label_keys_override: Option>, } impl InstantVectorElement { pub fn new(labels: KeyByLabelValues, value: f64) -> Self { - Self { labels, value } + Self { + labels, + value, + label_keys_override: None, + } + } + + /// Attach a per-element label-key override (see field doc on + /// `InstantVectorElement::label_keys_override`). + pub fn with_label_keys_override(mut self, keys: Vec) -> Self { + self.label_keys_override = Some(keys); + self } } diff --git a/data_plane/src/utils/http.rs b/data_plane/src/utils/http.rs index fad752b6..fe4a4d1a 100644 --- a/data_plane/src/utils/http.rs +++ b/data_plane/src/utils/http.rs @@ -135,10 +135,21 @@ pub fn convert_query_result_to_prometheus( let timestamp = instant_vector.timestamp as f64 / 1000.0; for element in &instant_vector.values { - // zip over query_output_labels.keys and element.labels.labels and collect into metric_map - let mut metric_map = HashMap::new(); - for (key, label) in query_output_labels - .labels + // Build metric labels object. Per-element override + // (`element.label_keys_override`) wins when the + // adapter knows the keys at materialization time — + // e.g. ASAP-tier `topk(...)` synthesizes an `"item"` + // key that's not in the query's group-by clause, so + // the outer `query_output_labels` doesn't carry it. + // Falls back to the query-scoped key list for everyone + // else. Mirrors the matrix branch in + // `convert_range_result_to_prometheus`. + let mut metric_map: HashMap<&String, &String> = HashMap::new(); + let effective_keys: &[String] = element + .label_keys_override + .as_deref() + .unwrap_or(&query_output_labels.labels); + for (key, label) in effective_keys .iter() .zip(element.labels.labels.iter()) { diff --git a/data_plane/tests/e2e_controller_plans_and_backend_serves.rs b/data_plane/tests/e2e_controller_plans_and_backend_serves.rs index bca83233..61455c0e 100644 --- a/data_plane/tests/e2e_controller_plans_and_backend_serves.rs +++ b/data_plane/tests/e2e_controller_plans_and_backend_serves.rs @@ -1677,15 +1677,11 @@ async fn controller_plan_to_query_full_roundtrip_cms_with_heap_topk() { "topk(...) on CmsWithHeap must succeed end-to-end. Response:\n{}", serde_json::to_string_pretty(&response).unwrap_or_default() ); - // Top-1 should be `gamma` (count=200). The reducer keys each - // top-k item by `item: ` in the series labels, but the - // wire-format `InstantVectorElement` adapter currently drops - // per-element labels (`label_keys_override` only exists on - // `RangeVectorElement`); confirmed by the response carrying - // `"metric": {}` on every element. Until that adapter gap is - // closed, assert the strongest invariants the wire-format DOES - // surface: `topk(3)` returned at least one series, the values - // include `gamma`'s count (200), and we got at most 3 results. + // Top-1 must be `gamma` (count=200), surfaced via the + // `item: ` synthesized label on each top-k series. The + // `InstantVectorElement::label_keys_override` field (added + // alongside this assertion's tightening) carries the synthesized + // key through the Prometheus adapter. let result = &response["data"]["result"]; let arr = result .as_array() @@ -1704,6 +1700,21 @@ async fn controller_plan_to_query_full_roundtrip_cms_with_heap_topk() { }) .collect(); values.sort_by(|a, b| b.partial_cmp(a).unwrap_or(std::cmp::Ordering::Equal)); + let mut found_gamma = false; + for elem in arr { + if let Some(item) = elem["metric"]["item"].as_str() { + if item == "gamma" { + found_gamma = true; + break; + } + } + } + assert!( + found_gamma, + "topk(3) must surface `gamma` via the `item` label on at least \ + one series. Response:\n{}", + serde_json::to_string_pretty(&response).unwrap_or_default() + ); assert!( values.first().map(|v| (v - 200.0).abs() < 1.0).unwrap_or(false), "topk(3) on heap-bearing CMS must surface `gamma`'s count (200) as \