From 8b4d6790927b662ec6cc119ba286b6e2cf7246c7 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 07:23:08 +0000 Subject: [PATCH 1/2] refactor: compute current-series readouts as Planner physical programs Storage computed current-series Sum, Count, Average, Quantile and TopK itself from cached per-group arrays; only TopK was a Planner program over the population snapshot. Populations are now selected over roots typed with the complete series identity, so Planner compiles every ReadPopulation. Storage keeps the population and returns its members; the entry's retained physical program computes the readout. The readout field, the storage readout caches and the cache-builds metric are removed. `without` populations are not selected (Planner cannot yet project them from the series identity); those queries run exactly. Co-Authored-By: Claude Opus 5.5 --- control_plane/src/physical/compiler.rs | 7 +- .../src/physical/maintained_population.rs | 71 ++-- control_plane/src/physical/workload_cost.rs | 50 +-- .../src/query_plan/current_series.rs | 16 - crates/asap_types/src/query_plan/native.rs | 27 +- .../asap_types/src/query_plan/query_time.rs | 23 +- data_plane/src/drivers/query/servers/http.rs | 13 +- .../query_engines/asap_query_engine/engine.rs | 2 - .../logical_dag/native_values.rs | 14 +- .../asap_query_engine/request_tests.rs | 1 - .../sketch_db/current_series.rs | 310 +++--------------- .../tests/support/current_series_process.rs | 4 - .../current-series-aggregations.md | 25 +- 13 files changed, 167 insertions(+), 396 deletions(-) diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 4be16cfd4..09ff1c1e5 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -2436,7 +2436,7 @@ impl DeploymentPlanCompiler { if frontend == QueryFrontend::MetricsQl { entry.language = crate::query_plan::QueryLanguage::MetricsQl; } - super::maintained_population::install_native_topk( + super::maintained_population::install_population_readout( &mut entry, query .retained_physical()? @@ -4867,6 +4867,11 @@ pub(crate) mod tests { } } assert_eq!(populations.len(), 1); + // Storage returns the members; every readout is the entry's Planner program. + for entry in plan.query_plan.entries.values() { + assert!(entry.population_snapshot().is_some(), "{entry:?}"); + entry.recover_population_physical_dag().unwrap(); + } let installed = serde_json::to_string(&plan.precompute_plan.executable_dags).unwrap(); assert!( !installed.contains("MaintainPopulation"), diff --git a/control_plane/src/physical/maintained_population.rs b/control_plane/src/physical/maintained_population.rs index c35e1c717..7606ba24d 100644 --- a/control_plane/src/physical/maintained_population.rs +++ b/control_plane/src/physical/maintained_population.rs @@ -4,17 +4,17 @@ use super::compiler::QueryCompilationInput; use super::compiler::{CompileError, PhysicalCompilationRequest}; use asap_types::physical_plan_codec::PhysicalPlanCodec; use asap_types::query_plan::{ - current_series::{SeriesPopulation, SeriesReadout}, + current_series::SeriesPopulation, query_time::{Grouping, LabelMatch, LabelMatcher, QueryTimeOperator}, }; use planner_types::post_asap::{ maintained_population::*, SummaryExpr, SummaryNode, ValueOperation, }; -fn selected(node: &SummaryNode) -> Option<(MaintainedPopulation, PopulationReadout)> { +fn selected(node: &SummaryNode) -> Option { if let SummaryExpr::ValueOperation { child, - operation: ValueOperation::ReadPopulation { readout }, + operation: ValueOperation::ReadPopulation { .. }, .. } = &node.expr { @@ -23,13 +23,13 @@ fn selected(node: &SummaryNode) -> Option<(MaintainedPopulation, PopulationReado .. } = &child.expr { - return Some((population.clone(), readout.clone())); + return Some(population.clone()); } } // The source remains a maintained population when Planner places a heap, // projection and ranking above it. Backend binds that source only. let SummaryExpr::ValueOperation { - operation: ValueOperation::Limit { n, offset: 0, .. }, + operation: ValueOperation::Limit { offset: 0, .. }, .. } = &node.expr else { @@ -54,12 +54,14 @@ fn selected(node: &SummaryNode) -> Option<(MaintainedPopulation, PopulationReado &std::rc::Rc::new(node.clone()), ) .ok()?; - Some(((*population).clone(), PopulationReadout::TopK { k: *n })) + Some((*population).clone()) } +/// A current-series population Planner can read out. Planner cannot yet +/// project `without` groups from the series identity. pub(super) fn supported_node(node: &SummaryNode) -> bool { - selected(node).is_some_and(|(population, _)| { - matches!(population.input, PopulationInput::CurrentSeries(_)) + selected(node).is_some_and(|population| { + matches!(&population.input, PopulationInput::CurrentSeries(spec) if !spec.without) }) } @@ -83,9 +85,7 @@ pub(super) fn operators( let populations: std::collections::BTreeSet<_> = selected .iter() .flatten() - .map(|(population, _)| { - serde_json::to_string(population).expect("typed population serializes") - }) + .map(|population| serde_json::to_string(population).expect("typed population serializes")) .collect(); let max_bytes = request .retained_summary_memory_budget_bytes @@ -93,7 +93,7 @@ pub(super) fn operators( .min(1_073_741_824) / populations.len().max(1) as u64; selected.into_iter().zip(&request.queries).map(|(selected, query)| { - let Some((spec, readout)) = selected else { return Ok(None); }; + let Some(spec) = selected else { return Ok(None); }; let PopulationInput::CurrentSeries(input) = &spec.input else { return Err(CompileError::Query { query_id: query.query_id.clone(), reason: "maintained table-row populations require a row-update executor; remote-write current-series state is incompatible".into() }); }; @@ -131,17 +131,7 @@ pub(super) fn operators( .min(input.lookback_ms), }; population.validate()?; - let readout = match &readout { - PopulationReadout::Quantile { q } => SeriesReadout::Quantile { q: *q }, - PopulationReadout::TopK { k } => SeriesReadout::TopK { k: *k as u64 }, - PopulationReadout::Sum => SeriesReadout::Sum, - PopulationReadout::Count => SeriesReadout::Count, - PopulationReadout::Average => SeriesReadout::Average, - }; - Ok(Some(QueryTimeOperator::CurrentSeries { - population, - readout, - })) + Ok(Some(QueryTimeOperator::CurrentSeries { population })) }).collect() } @@ -158,30 +148,25 @@ pub(super) fn operator( Ok(operators(request)?.remove(index)) } -/// The maintained population is a deployment source; ranking is compiled by -/// Planner before this candidate is priced or installed. -pub(super) fn install_native_topk( +/// The maintained population is a deployment source; its readout is the +/// Planner program compiled before this candidate is priced or installed. +pub(super) fn install_population_readout( entry: &mut asap_types::query_plan::QueryPlanEntry, compiled: Option<&asap_physical_operators::physical_planner::CompiledPhysicalDag>, ) -> Result<(), CompileError> { use asap_types::query_plan::QueryPlanNode; - let Some(QueryPlanNode::Logical { - operator: - QueryTimeOperator::CurrentSeries { - population, - readout: SeriesReadout::TopK { .. }, - }, - .. - }) = entry.nodes.get(&entry.root) - else { - return Ok(()); - }; - if population.grouping.without { + if !matches!( + entry.nodes.get(&entry.root), + Some(QueryPlanNode::Logical { + operator: QueryTimeOperator::CurrentSeries { .. }, + .. + }) + ) { return Ok(()); } let compiled = compiled.ok_or_else(|| CompileError::Query { query_id: entry.query_id.clone(), - reason: "selected TopK candidate has no retained physical DAG".into(), + reason: "selected population readout has no retained physical DAG".into(), })?; let encoded = compiled.encode().map_err(|error| CompileError::Query { query_id: entry.query_id.clone(), @@ -191,14 +176,6 @@ pub(super) fn install_native_topk( serde_json::from_slice(&encoded) .map_err(|error| CompileError::Snapshot(error.to_string()))?, ); - let Some(QueryPlanNode::Logical { - operator: QueryTimeOperator::CurrentSeries { readout, .. }, - .. - }) = entry.nodes.get_mut(&entry.root) - else { - unreachable!() - }; - *readout = SeriesReadout::Snapshot; entry.recover_population_physical_dag()?; Ok(()) } diff --git a/control_plane/src/physical/workload_cost.rs b/control_plane/src/physical/workload_cost.rs index 56d3b5d8a..db0a9966c 100644 --- a/control_plane/src/physical/workload_cost.rs +++ b/control_plane/src/physical/workload_cost.rs @@ -834,13 +834,24 @@ fn enumerate_frontier_candidates( _ => unreachable!("native candidate retains canonical roots"), }) .collect(); - let strategy = - asap_aware_mapping::maintained_population::MaintainedPopulationStrategy::new(&roots); - let maintained_roots: Vec<_> = roots + // Populations are selected over roots typed with the complete series + // identity, so Planner can compile every readout over the snapshot rows. + let typed = roots .iter() .map(|root| { - strategy - .candidate(root) + asap_physical_operators::physical_planner::promql_rows::with_series_identity(root) + .map(std::rc::Rc::new) + .ok() + }) + .collect::>(); + let strategy = asap_aware_mapping::maintained_population::MaintainedPopulationStrategy::new( + &typed.iter().flatten().cloned().collect::>(), + ); + let maintained_roots: Vec<_> = typed + .iter() + .map(|root| { + root.as_ref() + .and_then(|root| strategy.candidate(root)) .filter(|node| super::maintained_population::supported_node(node)) }) .collect(); @@ -1050,20 +1061,21 @@ mod tests { queries.truncate(1); queries[0].query = planner_types::workload::Query("count by(job)(m)".into()); let plan = with_unit_quotes(input).compile_promql().unwrap(); - assert!(plan - .query_plan - .entries - .values() - .all(|entry| entry.nodes.values().any(|node| matches!( - node, - crate::query_plan::QueryPlanNode::Logical { - operator: crate::query_plan::query_time::QueryTimeOperator::CurrentSeries { - readout: asap_types::query_plan::current_series::SeriesReadout::Count, - .. - }, - .. - } - )))); + // The population supplies members; Planner counts them per job. + for entry in plan.query_plan.entries.values() { + assert!(entry.population_snapshot().is_some(), "{entry:?}"); + let program = entry.recover_population_physical_dag().unwrap(); + let output = program.output_contract(program.roots()[0]).unwrap(); + assert_eq!( + output + .schema + .fields + .iter() + .map(|field| field.name.as_str()) + .collect::>(), + ["job", "count"] + ); + } } fn fixture() -> BackendLocalPlanningInput { diff --git a/crates/asap_types/src/query_plan/current_series.rs b/crates/asap_types/src/query_plan/current_series.rs index 55fb01f67..1f8ffd1e8 100644 --- a/crates/asap_types/src/query_plan/current_series.rs +++ b/crates/asap_types/src/query_plan/current_series.rs @@ -46,22 +46,6 @@ impl SeriesPopulation { } } -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] -#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)] -pub enum SeriesReadout { - /// All eligible members for a Planner-compiled physical readout. - Snapshot, - Quantile { - q: f64, - }, - TopK { - k: u64, - }, - Sum, - Count, - Average, -} - #[cfg(test)] mod tests { use super::*; diff --git a/crates/asap_types/src/query_plan/native.rs b/crates/asap_types/src/query_plan/native.rs index 42b10e3da..dc8bc9029 100644 --- a/crates/asap_types/src/query_plan/native.rs +++ b/crates/asap_types/src/query_plan/native.rs @@ -11,11 +11,7 @@ impl QueryPlanEntry { } match self.nodes.get(&self.root) { Some(QueryPlanNode::Logical { - operator: - query_time::QueryTimeOperator::CurrentSeries { - population, - readout: current_series::SeriesReadout::Snapshot, - }, + operator: query_time::QueryTimeOperator::CurrentSeries { population }, inputs, }) if inputs.is_empty() => Some(population), _ => None, @@ -49,12 +45,27 @@ impl QueryPlanEntry { let output = dag .output_contract(*root) .map_err(|error| QueryPlanError::Invalid(error.to_string()))?; - if input.schema != output.schema { + use planner_types::{post_asap::SummaryFamilyType, pre_asap::DataType}; + // A ranking returns complete source rows; an aggregate returns one + // value per group of the population's grouping labels. + let numeric = |field: &planner_types::post_asap::SummaryField| { + matches!( + field.dtype, + SummaryFamilyType::Plain(DataType::Float64 | DataType::Int64) + ) + }; + let grouped_value = output.schema.time_index.is_none() + && output.schema.fields.iter().filter(|f| numeric(f)).count() == 1 + && output.schema.fields.iter().all(|field| { + numeric(field) + || (field.dtype == SummaryFamilyType::Plain(DataType::Utf8) + && population.grouping.labels.contains(&field.name)) + }); + if input.schema != output.schema && !grouped_value { return Err(invalid( - "population ranking must preserve complete source rows", + "population readout must return source rows or one value per group", )); } - use planner_types::{post_asap::SummaryFamilyType, pre_asap::DataType}; let fields = &input.schema.fields; let column = |name: &str, dtype: DataType| { fields.iter().any(|field| { diff --git a/crates/asap_types/src/query_plan/query_time.rs b/crates/asap_types/src/query_plan/query_time.rs index a67942f11..07542a740 100644 --- a/crates/asap_types/src/query_plan/query_time.rs +++ b/crates/asap_types/src/query_plan/query_time.rs @@ -10,10 +10,10 @@ fn invalid(message: impl Into) -> QueryPlanError { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] #[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)] pub enum QueryTimeOperator { - /// Readout over a bounded current-value population maintained at ingest. + /// Every member of a bounded current-value population maintained at + /// ingest; the entry's Planner program computes the readout. CurrentSeries { population: super::current_series::SeriesPopulation, - readout: super::current_series::SeriesReadout, }, /// A maximal exact scalar/vector subtree evaluated by Prometheus. ExactSubquery { query: String }, @@ -51,25 +51,8 @@ pub enum LabelMatch { } impl QueryTimeOperator { pub fn validate(&self, inputs: usize) -> Result<(), QueryPlanError> { - if let Self::CurrentSeries { - population, - readout, - } = self - { + if let Self::CurrentSeries { population } = self { population.validate()?; - match readout { - super::current_series::SeriesReadout::Quantile { q } - if !q.is_finite() || !population.quantiles => - { - return Err(invalid( - "quantile readout requires finite q and a quantile population", - )) - } - super::current_series::SeriesReadout::TopK { k } if *k > population.max_k => { - return Err(invalid("TopK readout exceeds shared population capacity")) - } - _ => {} - } } let expected = match self { Self::Scan { .. } | Self::ExactSubquery { .. } | Self::CurrentSeries { .. } => 0, diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index c48391baa..7a0c94dc1 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -2053,15 +2053,18 @@ async fn handle_metrics(State(state): State) -> impl IntoResponse { let mut buffer = Vec::new(); prometheus::Encoder::encode(&encoder, &metric_families, &mut buffer) .unwrap_or_else(|e| tracing::error!("Failed to encode metrics: {}", e)); - let (populations, builds) = state + let populations = state .summary_store .current_series .lock() .expect("current-series state poisoned") - .stats(); - buffer.extend_from_slice(format!( - "# TYPE asap_current_series_populations gauge\nasap_current_series_populations {populations}\n# TYPE asap_current_series_cache_builds_total counter\nasap_current_series_cache_builds_total {builds}\n" - ).as_bytes()); + .population_count(); + buffer.extend_from_slice( + format!( + "# TYPE asap_current_series_populations gauge\nasap_current_series_populations {populations}\n" + ) + .as_bytes(), + ); if let Some(receiver) = state.remote_write.as_ref() { use std::sync::atomic::Ordering; let stats = receiver.stats(); 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 824081016..46f798d5f 100644 --- a/data_plane/src/query_engines/asap_query_engine/engine.rs +++ b/data_plane/src/query_engines/asap_query_engine/engine.rs @@ -549,7 +549,6 @@ impl ASAPQueryEngine { operator: asap_types::query_plan::query_time::QueryTimeOperator::CurrentSeries { population, - readout, }, .. }) = entry.nodes.get(&root) @@ -567,7 +566,6 @@ impl ASAPQueryEngine { physical.query_plan.plan_version, ), population, - readout, evaluation_ms, ) .map_err(|error| EngineError::capability_miss("current_series", error))?; diff --git a/data_plane/src/query_engines/asap_query_engine/logical_dag/native_values.rs b/data_plane/src/query_engines/asap_query_engine/logical_dag/native_values.rs index e62a0409a..f0b37c3f3 100644 --- a/data_plane/src/query_engines/asap_query_engine/logical_dag/native_values.rs +++ b/data_plane/src/query_engines/asap_query_engine/logical_dag/native_values.rs @@ -871,11 +871,19 @@ fn execute_batches( .schema() .fields .iter() - .position(|field| field.dtype == SummaryFamilyType::Plain(DataType::Float64)) + .position(|field| { + matches!( + field.dtype, + SummaryFamilyType::Plain(DataType::Float64 | DataType::Int64) + ) + }) .ok_or_else(|| miss("physical output loses sample value"))?; for row in batch.rows() { - let Value::Float64(sample) = &row[value] else { - return Err(miss("invalid physical result value")); + let sample = &match row[value] { + Value::Float64(sample) => sample, + // A count is exact as a PromQL sample up to 2^53. + Value::Int64(count) if count.unsigned_abs() <= 1 << 53 => count as f64, + _ => return Err(miss("invalid physical result value")), }; let labels = if let Some(identity) = identity { let Value::Utf8(encoded) = &row[identity] else { diff --git a/data_plane/src/query_engines/asap_query_engine/request_tests.rs b/data_plane/src/query_engines/asap_query_engine/request_tests.rs index 2ae8e83c6..3f876c26a 100644 --- a/data_plane/src/query_engines/asap_query_engine/request_tests.rs +++ b/data_plane/src/query_engines/asap_query_engine/request_tests.rs @@ -127,7 +127,6 @@ fn mixed_bound_outputs_rejected_before_any_source_is_read() { max_k: 3, quantiles: false, }, - readout: current_series::SeriesReadout::Sum, }, inputs: vec![], }, diff --git a/data_plane/src/storage_engines/sketch_db/current_series.rs b/data_plane/src/storage_engines/sketch_db/current_series.rs index f9bd10728..e5d415023 100644 --- a/data_plane/src/storage_engines/sketch_db/current_series.rs +++ b/data_plane/src/storage_engines/sketch_db/current_series.rs @@ -1,7 +1,7 @@ //! Bounded current-value state. Never pools a series' old samples into a quantile. use crate::drivers::ingest::prometheus_remote_write::CanonicalSample; use asap_types::query_plan::{ - current_series::{SeriesPopulation, SeriesReadout}, + current_series::SeriesPopulation, query_time::{LabelMatch, QueryTimeOperator}, QueryPlan, QueryPlanNode, }; @@ -42,7 +42,6 @@ struct Member { #[derive(Default, Clone)] struct Group { ordered: BTreeSet, - cached: Option<(Vec, Vector, f64, f64)>, } #[derive(Clone)] struct Population { @@ -54,7 +53,6 @@ struct Population { /// Input timestamp at which the budget blew, cleared once the lookback /// window has moved entirely past it. `None` means the population serves. unavailable: Option, - cache_builds: u64, last_read: i64, matchers: Vec, } @@ -83,7 +81,6 @@ impl Population { groups: BTreeMap::new(), bytes: 0, unavailable: None, - cache_builds: 0, last_read: i64::MIN, matchers, }) @@ -98,7 +95,6 @@ impl Population { value, labels: labels.clone(), }); - group.cached = None; if group.ordered.is_empty() { self.groups.remove(&old.group); } @@ -190,94 +186,16 @@ impl Population { if let Some(value) = sample.value { let state = self.groups.entry(group).or_default(); state.ordered.insert(Ranked { value, labels }); - state.cached = None; } } - fn read(&mut self, readout: &SeriesReadout) -> Vector { - if matches!(readout, SeriesReadout::Snapshot) { - return self - .groups - .values() - .flat_map(|group| group.ordered.iter()) - .map(|member| (member.labels.clone(), member.value)) - .collect(); - } - let mut result = vec![]; - for (labels, group) in &mut self.groups { - if group.cached.is_none() { - let values = if self.definition.quantiles { - group.ordered.iter().map(|r| r.value).collect() - } else { - vec![] - }; - let top = group - .ordered - .iter() - .rev() - .take(self.definition.max_k as usize) - .map(|r| (r.labels.clone(), r.value)) - .collect(); - let sum = compensated_sum(group.ordered.iter().map(|r| r.value)); - let count = group.ordered.len() as f64; - let average = if sum.is_finite() { - sum / count - } else { - compensated_sum(group.ordered.iter().map(|r| r.value / count)) - }; - group.cached = Some((values, top, sum, average)); - self.cache_builds += 1; - } - let (values, top, sum, average) = group.cached.as_ref().unwrap(); - match readout { - SeriesReadout::Snapshot => unreachable!("snapshot returned above"), - SeriesReadout::Quantile { q } => { - // `values` is only populated for a quantile-carrying population. - // `QueryTimeOperator::validate` rejects the mismatched pairing at - // install, so this is defensive: answer like Prometheus does for - // an empty group rather than underflow `values.len() - 1` while - // holding the lock every remote-write batch waits on. - let value = if values.is_empty() { - f64::NAN - } else if *q < 0. { - f64::NEG_INFINITY - } else if *q > 1. { - f64::INFINITY - } else { - let rank = q * (values.len() - 1) as f64; - let lo = rank.floor() as usize; - let hi = (lo + 1).min(values.len() - 1); - let weight = rank - lo as f64; - values[lo] * (1. - weight) + values[hi] * weight - }; - result.push((labels.clone(), value)); - } - SeriesReadout::TopK { k } => result.extend(top.iter().take(*k as usize).cloned()), - SeriesReadout::Sum => result.push((labels.clone(), *sum)), - SeriesReadout::Count => result.push((labels.clone(), group.ordered.len() as f64)), - SeriesReadout::Average => result.push((labels.clone(), *average)), - } - } - result - } -} - -// Rebuild shared statistics after replacement/expiry, avoiding subtraction drift. -fn compensated_sum(values: impl Iterator) -> f64 { - let (mut sum, mut correction) = (0.0_f64, 0.0); - for value in values { - let next = sum + value; - if next.is_finite() { - correction += if sum.abs() >= value.abs() { - (sum - next) + value - } else { - (value - next) + sum - }; - } else { - correction = 0.0; - } - sum = next; + /// Every current member with a value, group by group. + fn snapshot(&self) -> Vector { + self.groups + .values() + .flat_map(|group| group.ordered.iter()) + .map(|member| (member.labels.clone(), member.value)) + .collect() } - sum + correction } #[derive(Default)] @@ -404,7 +322,6 @@ impl CurrentSeriesStore { &mut self, generation: (u64, u64), definition: &SeriesPopulation, - readout: &SeriesReadout, at: u64, ) -> Result { if self.generation != Some(generation) { @@ -436,7 +353,7 @@ impl CurrentSeriesStore { } let mut population = saved.clone(); population.expire(at.saturating_sub(definition.lookback_ms as i64)); - return Ok(population.read(readout)); + return Ok(population.snapshot()); } if at > watermark.saturating_add(definition.max_input_lag_ms as i64) { return Err("current-series input is behind evaluation time".into()); @@ -459,13 +376,10 @@ impl CurrentSeriesStore { } population.last_read = at; population.expire(at.saturating_sub(definition.lookback_ms as i64)); - Ok(population.read(readout)) + Ok(population.snapshot()) } - pub fn stats(&self) -> (usize, u64) { - ( - self.populations.len(), - self.populations.values().map(|p| p.cache_builds).sum(), - ) + pub fn population_count(&self) -> usize { + self.populations.len() } } @@ -508,14 +422,10 @@ mod tests { let mut store = CurrentSeriesStore::default(); store.ingest(&plan, &[sample("old", "api", 0, Some(10.))]); store.ingest(&plan, &[sample("new", "api", 500, Some(3.))]); - let values = store - .read((7, 1), &population, &SeriesReadout::Sum, 1_000) - .unwrap(); + let values = store.read((7, 1), &population, 1_000).unwrap(); assert_eq!(values.len(), 1); assert_eq!(values[0].1, 3.); - let values = store - .read((7, 1), &population, &SeriesReadout::Sum, 1_500) - .unwrap(); + let values = store.read((7, 1), &population, 1_500).unwrap(); assert!(values.is_empty()); } // Historical reads use the state at that timestamp, never the latest values. @@ -536,23 +446,9 @@ mod tests { sample("one", "api", 3_000, Some(4.)), ], ); - assert_eq!( - store - .read((7, 1), &population, &SeriesReadout::Sum, 1_000) - .unwrap()[0] - .1, - 2. - ); - assert_eq!( - store - .read((7, 1), &population, &SeriesReadout::Sum, 3_000) - .unwrap()[0] - .1, - 4. - ); - assert!(store - .read((7, 1), &population, &SeriesReadout::Sum, 999) - .is_err()); + assert_eq!(store.read((7, 1), &population, 1_000).unwrap()[0].1, 2.); + assert_eq!(store.read((7, 1), &population, 3_000).unwrap()[0].1, 4.); + assert!(store.read((7, 1), &population, 999).is_err()); assert!( store.history_bytes[&population.key()] + store.populations[&population.key()].bytes <= population.max_bytes @@ -577,9 +473,7 @@ mod tests { ], ); store.ingest(&installed, &[sample("two", "api", 1_500, Some(20.))]); - assert!(store - .read((7, 1), &population, &SeriesReadout::Sum, 1_000) - .is_err()); + assert!(store.read((7, 1), &population, 1_000).is_err()); let mut bounded = population.clone(); bounded.max_bytes = 2_500; let mut store = CurrentSeriesStore::default(); @@ -591,9 +485,7 @@ mod tests { sample("one", "api", 2_000, Some(3.)), ], ); - assert!(store - .read((7, 1), &bounded, &SeriesReadout::Sum, 1_000) - .is_err()); + assert!(store.read((7, 1), &bounded, 1_000).is_err()); assert!( store.history_bytes[&bounded.key()] + store.populations[&bounded.key()].bytes <= bounded.max_bytes @@ -618,7 +510,6 @@ mod tests { QueryPlanNode::Logical { operator: QueryTimeOperator::CurrentSeries { population: p.clone(), - readout: SeriesReadout::Quantile { q: 0.5 }, }, inputs: vec![], }, @@ -670,99 +561,40 @@ mod tests { ); } } - // Equal sample values still represent two series; replacements and stale markers retract them. - #[test] - fn sum_count_average_follow_current_series_membership() { - let p = definition(); - let plan = plan(&p); - let mut store = CurrentSeriesStore::default(); - warm(&mut store, &plan); - for (at, samples, expected) in [ - ( - 301_000, - vec![sample("y", "api", 301_000, Some(1.))], - [7., 3., 7. / 3.], - ), - ( - 302_000, - vec![sample("z", "api", 302_000, None)], - [2., 2., 1.], - ), - ] { - store.ingest(&plan, &samples); - for (readout, truth) in [ - SeriesReadout::Sum, - SeriesReadout::Count, - SeriesReadout::Average, - ] + fn members(values: Vector) -> Vec<(String, f64)> { + let mut members = values .into_iter() - .zip(expected) - { - let values = store.read((7, 1), &p, &readout, at).unwrap(); - assert!( - (values[0].1 - truth).abs() < 1e-12, - "{readout:?}: {values:?}" - ); - } - } + .map(|(labels, value)| (labels["pod"].clone(), value)) + .collect::>(); + members.sort_by(|a, b| a.0.cmp(&b.0)); + members } - /// Four quantiles reuse one distribution, and smaller k reads the shared maximum-k prefix. + // Equal sample values still represent two series; replacements and stale + // markers retract them, and out-of-order values cannot resurrect a series. #[test] - fn quantiles_and_topk_share_state_and_promote_after_updates_and_staleness() { + fn snapshot_follows_current_series_membership() { let p = definition(); let plan = plan(&p); let mut store = CurrentSeriesStore::default(); warm(&mut store, &plan); - for (q, expected) in [(0.5, 5.), (0.9, 8.2), (0.95, 8.6), (0.99, 8.92)] { - let result = store - .read((7, 1), &p, &SeriesReadout::Quantile { q }, 300_000) - .unwrap(); - assert!((result[0].1 - expected).abs() < 1e-10); - assert_eq!(result[1].1, 50.); - } - for percentile in 1..100 { - let q = percentile as f64 / 100.; - let result = store - .read((7, 1), &p, &SeriesReadout::Quantile { q }, 300_000) - .unwrap(); - assert!((result[0].1 - (1. + 8. * q)).abs() < 1e-10); - } - let small = store - .read((7, 1), &p, &SeriesReadout::TopK { k: 1 }, 300_000) - .unwrap(); - let big = store - .read((7, 1), &p, &SeriesReadout::TopK { k: 3 }, 300_000) - .unwrap(); - assert_eq!(small[0].0["pod"], "y"); - assert_eq!(big[0], small[0]); - assert_eq!(store.stats(), (1, 2)); - store.ingest(&plan, &[sample("y", "api", 301_000, Some(-5.))]); + store.ingest(&plan, &[sample("y", "api", 301_000, Some(1.))]); assert_eq!( - store - .read((7, 1), &p, &SeriesReadout::TopK { k: 1 }, 301_000) - .unwrap()[0] - .0["pod"], - "z" + members(store.read((7, 1), &p, 301_000).unwrap()), + [("w", 50.), ("x", 1.), ("y", 1.), ("z", 5.)].map(|(pod, value)| (pod.into(), value)) ); store.ingest(&plan, &[sample("z", "api", 302_000, None)]); - assert_eq!( - store - .read((7, 1), &p, &SeriesReadout::TopK { k: 1 }, 302_000) - .unwrap()[0] - .0["pod"], - "x" - ); - // Out-of-order old values must not resurrect the stale series. store.ingest(&plan, &[sample("z", "api", 301_000, Some(100.))]); assert_eq!( - store - .read((7, 1), &p, &SeriesReadout::TopK { k: 1 }, 302_000) - .unwrap()[0] - .0["pod"], - "x" + members(store.read((7, 1), &p, 302_000).unwrap()), + [("w", 50.), ("x", 1.), ("y", 1.)].map(|(pod, value)| (pod.into(), value)) ); + let rows = store.read((7, 1), &p, 302_000).unwrap(); + assert!(rows + .iter() + .all(|(labels, _)| labels["__name__"] == "a" && labels.contains_key("job"))); } + /// Cold state, gaps, old generations and historical timestamps cannot masquerade as complete populations. #[test] fn coverage_expiration_generation_and_capacity_fail_closed() { @@ -770,59 +602,32 @@ mod tests { let plan = plan(&p); let mut store = CurrentSeriesStore::default(); store.ingest(&plan, &[sample("x", "api", 0, Some(1.))]); - assert!(store - .read((7, 1), &p, &SeriesReadout::Quantile { q: 0.5 }, 0) - .is_err()); + assert!(store.read((7, 1), &p, 0).is_err()); warm(&mut store, &plan); - assert!(store - .read((7, 2), &p, &SeriesReadout::Quantile { q: 0.5 }, 300_000) - .is_err()); - assert!(store - .read((7, 1), &p, &SeriesReadout::Quantile { q: 0.5 }, 299_000) - .is_err()); + assert!(store.read((7, 2), &p, 300_000).is_err()); + assert!(store.read((7, 1), &p, 299_000).is_err()); for t in (360_000..=600_000).step_by(60_000) { store.ingest(&plan, &[sample("y", "api", t, Some(9.))]); } - assert_eq!( - store - .read((7, 1), &p, &SeriesReadout::Quantile { q: 0.5 }, 600_000) - .unwrap() - .len(), - 1 - ); // Prometheus 3.5 lookback is left-open - assert_eq!( - store - .read((7, 1), &p, &SeriesReadout::Quantile { q: 0.5 }, 600_001) - .unwrap() - .len(), - 1 - ); - assert!(store - .read((7, 1), &p, &SeriesReadout::Quantile { q: 0.5 }, 600_000) - .is_err()); + assert_eq!(store.read((7, 1), &p, 600_000).unwrap().len(), 1); // Prometheus 3.5 lookback is left-open + assert_eq!(store.read((7, 1), &p, 600_001).unwrap().len(), 1); + assert!(store.read((7, 1), &p, 600_000).is_err()); store.ingest(&plan, &[sample("y", "api", 900_000, Some(9.))]); - assert!(store - .read((7, 1), &p, &SeriesReadout::Quantile { q: 0.5 }, 900_000) - .is_err()); + assert!(store.read((7, 1), &p, 900_000).is_err()); let mut bounded = p.clone(); bounded.max_series = 3; let plan = super::tests::plan(&bounded); let mut store = CurrentSeriesStore::default(); warm(&mut store, &plan); assert!(store - .read( - (7, 1), - &bounded, - &SeriesReadout::Quantile { q: 0.5 }, - 300_000 - ) + .read((7, 1), &bounded, 300_000) .unwrap_err() .contains("budget")); } // Prometheus 3.5 selectors exclude samples exactly at evaluation - lookback. #[test] - fn lookback_left_boundary_expires_members_for_all_shared_readouts() { + fn lookback_left_boundary_expires_members() { let p = definition(); let plan = plan(&p); let mut store = CurrentSeriesStore::default(); @@ -830,22 +635,9 @@ mod tests { for t in (360_000..=600_000).step_by(60_000) { store.ingest(&plan, &[sample("y", "api", t, Some(9.))]); } - for (readout, expected) in [ - (SeriesReadout::Count, 1.0), - (SeriesReadout::Sum, 9.0), - (SeriesReadout::Average, 9.0), - (SeriesReadout::Quantile { q: 0.5 }, 9.0), - ] { - let rows = store.read((7, 1), &p, &readout, 600_000).unwrap(); - assert_eq!( - rows, - vec![(BTreeMap::from([("job".into(), "api".into())]), expected)] - ); - } - let rows = store - .read((7, 1), &p, &SeriesReadout::TopK { k: 3 }, 600_000) - .unwrap(); - assert_eq!(rows.len(), 1); - assert_eq!(rows[0].0["pod"], "y"); + assert_eq!( + members(store.read((7, 1), &p, 600_000).unwrap()), + vec![("y".to_string(), 9.)] + ); } } diff --git a/data_plane/tests/support/current_series_process.rs b/data_plane/tests/support/current_series_process.rs index 86bb390b5..f94ce6db8 100644 --- a/data_plane/tests/support/current_series_process.rs +++ b/data_plane/tests/support/current_series_process.rs @@ -297,10 +297,6 @@ async fn current_series_quantiles_topk_share_and_replace_values() { metrics.contains("asap_current_series_populations 2\n"), "{metrics}" ); - assert!( - metrics.contains("asap_current_series_cache_builds_total 3\n"), - "{metrics}" - ); assert_eq!(calls.load(std::sync::atomic::Ordering::Relaxed), 0); // A decreasing value and a stale marker must promote a formerly excluded series. for (pod, value, offset, expected) in [ diff --git a/docs/developer_docs/query-engine/current-series-aggregations.md b/docs/developer_docs/query-engine/current-series-aggregations.md index 77fda0b7f..1c1eced54 100644 --- a/docs/developer_docs/query-engine/current-series-aggregations.md +++ b/docs/developer_docs/query-engine/current-series-aggregations.md @@ -1,13 +1,15 @@ # Shared current-series quantiles and TopK Backend-local PromQL workload compilation can export a maintained current-value -alternative for `quantile(q, metric)` and `topk(k, metric)`, including `by` and -`without` grouping and selector label matchers. Parameters must be finite scalar +alternative for `quantile(q, metric)`, `topk(k, metric)`, `sum`, `count` and +`avg`, with `by` grouping and selector label matchers. `without` grouping is +forwarded to the exact engine until Planner can project it from the series +identity. Parameters must be finite scalar literals. Selector offsets, `@`, nested input expressions and MetricsQL use the existing alternatives; they are not admitted by this implementation. ASAPPlanner owns this transformation through the opt-in `MaintainedPopulationStrategy` -over canonical IR. It emits `MaintainPopulation` at maintenance time and +over canonical IR typed with the complete series identity. It emits `MaintainPopulation` at maintenance time and `ReadPopulation` at read time, with source/filter/group identity, quantile consumers and the maximum requested k in its typed contract. Compatible producers are shared by Planner CSE. The backend consumes these nodes, binds resource and @@ -18,7 +20,10 @@ opt in only when they can implement and price this maintenance contract. For example, put p50, p90, p95, p99 and Top1/Top5 in one workload. Matching source, selector and grouping contracts produce one `CurrentSeries` population with `quantiles: true` and `max_k: 5`. Each registered query keeps its own readout. -Different sources, filters and groupings remain distinct. This state is exact: +Different sources, filters and groupings remain distinct. The backend stores +only the population and returns its current members; each readout (ranking, +quantile, sum, count, average) is the query's Planner-compiled physical program +over those members. This state is exact: it retains each series' current value in a shared ordered population, rather than inserting all historical observations into a quantile sketch. Its memory grows with series cardinality, even when only Top1 is requested. Those excluded series @@ -26,9 +31,7 @@ are required to promote the correct replacement when a winner decreases or expir Accepted Remote Write batches update the state atomically under its lock. A newer sample replaces the old value; stale markers remove the value. Older updates do -not resurrect a newer stale marker. Per-group readout arrays are shared until that -group changes. TopK-only populations cache just the largest registered k results; -quantile consumers also share an ordered value array. +not resurrect a newer stale marker. Deployment uses complete workload quotes. The manifest deduplicates population build, update, residency and retirement components across consumers and prices @@ -48,12 +51,12 @@ A native exact alternative remains available for cost selection and execution fa - State is in memory. Restart and generation replacement require warmup again; historical range queries continue to use native execution. - Populations divide the configured retained-summary memory budget and cap series - cardinality. Bounds include conservative space for labels, trees and caches. + cardinality. Bounds include conservative space for labels and trees. The existing Remote Write adapter accepts finite sample values and stale markers. -`/metrics` exposes `asap_current_series_populations` and -`asap_current_series_cache_builds_total` to verify reuse. These describe the active -in-memory population generation; they are not window-sketch materialization counts. +`/metrics` exposes `asap_current_series_populations` to verify sharing. It +describes the active in-memory population generation; it is not a window-sketch +materialization count. ## Validation From 9dfc40c39163522b83ad174caf57254ac2cb1f9a Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 07:49:50 +0000 Subject: [PATCH 2/2] fix: select only compilable populations; tighten readout decoding Review follow-ups: a population is selected only if Planner compiles its readout, so a failed compile runs the query exactly instead of failing the deployment. Grouped population outputs must carry exactly the grouping labels. An Int64 count is the sample only when it is the sole numeric column. The summation difference from Prometheus is documented. Co-Authored-By: Claude Opus 5.5 --- control_plane/src/physical/workload_cost.rs | 38 ++++++++++- crates/asap_types/src/query_plan/native.rs | 15 +++-- .../logical_dag/native_values.rs | 67 ++++++++++++++++--- .../sketch_db/current_series.rs | 2 +- .../current-series-aggregations.md | 5 +- 5 files changed, 109 insertions(+), 18 deletions(-) diff --git a/control_plane/src/physical/workload_cost.rs b/control_plane/src/physical/workload_cost.rs index db0a9966c..5a623f1d6 100644 --- a/control_plane/src/physical/workload_cost.rs +++ b/control_plane/src/physical/workload_cost.rs @@ -852,7 +852,12 @@ fn enumerate_frontier_candidates( .map(|root| { root.as_ref() .and_then(|root| strategy.candidate(root)) - .filter(|node| super::maintained_population::supported_node(node)) + .filter(|node| { + // Every readout must compile to the Planner program installed with it. + super::maintained_population::supported_node(node) + && asap_physical_operators::physical_planner::promql_rows::compile_current_series_readout(node) + .is_ok() + }) }) .collect(); if maintained_roots.iter().any(Option::is_some) { @@ -1048,6 +1053,37 @@ mod tests { }))); } + /// Planner cannot yet read out `without` groups, so such a population is + /// never selected, while the `by` form is. + #[test] + fn without_grouping_selects_no_population() { + for (query, supported) in [ + ("quantile without (pod) (0.5, m)", false), + ("quantile by (job) (0.5, m)", true), + ] { + let root = std::rc::Rc::new( + asap_physical_operators::physical_planner::promql_rows::with_series_identity( + &crate::query_parser::parse_query_expr_canonical( + query, + crate::types::AccuracyTarget::Exact, + ) + .unwrap(), + ) + .unwrap(), + ); + let strategy = + asap_aware_mapping::maintained_population::MaintainedPopulationStrategy::new( + std::slice::from_ref(&root), + ); + let candidate = strategy.candidate(&root).unwrap(); + assert_eq!( + super::super::maintained_population::supported_node(&candidate), + supported, + "{query}" + ); + } + } + /// Instant counts select current membership, never accumulated observations. #[test] fn local_grouped_count_has_a_bindable_candidate() { diff --git a/crates/asap_types/src/query_plan/native.rs b/crates/asap_types/src/query_plan/native.rs index dc8bc9029..6448b4002 100644 --- a/crates/asap_types/src/query_plan/native.rs +++ b/crates/asap_types/src/query_plan/native.rs @@ -54,12 +54,19 @@ impl QueryPlanEntry { SummaryFamilyType::Plain(DataType::Float64 | DataType::Int64) ) }; + let labels = output + .schema + .fields + .iter() + .filter(|field| !numeric(field)) + .map(|field| { + (field.dtype == SummaryFamilyType::Plain(DataType::Utf8)).then_some(&field.name) + }) + .collect::>>(); let grouped_value = output.schema.time_index.is_none() && output.schema.fields.iter().filter(|f| numeric(f)).count() == 1 - && output.schema.fields.iter().all(|field| { - numeric(field) - || (field.dtype == SummaryFamilyType::Plain(DataType::Utf8) - && population.grouping.labels.contains(&field.name)) + && labels.is_some_and(|labels| { + labels == population.grouping.labels.iter().collect::>() }); if input.schema != output.schema && !grouped_value { return Err(invalid( diff --git a/data_plane/src/query_engines/asap_query_engine/logical_dag/native_values.rs b/data_plane/src/query_engines/asap_query_engine/logical_dag/native_values.rs index f0b37c3f3..c6c711562 100644 --- a/data_plane/src/query_engines/asap_query_engine/logical_dag/native_values.rs +++ b/data_plane/src/query_engines/asap_query_engine/logical_dag/native_values.rs @@ -508,6 +508,43 @@ mod tests { ); } + // A population count is an Int64 sample only when it is the sole numeric + // column; a Float64 value always wins over an Int64 helper column. + #[test] + fn selected_program_reads_counts_and_prefers_float_values() { + for (fields, row, expected) in [ + ( + vec![("job", DataType::Utf8), ("count", DataType::Int64)], + vec![Value::Utf8("api".into()), Value::Int64(3)], + 3., + ), + ( + vec![("members", DataType::Int64), ("value", DataType::Float64)], + vec![Value::Int64(7), Value::Float64(2.5)], + 2.5, + ), + ] { + let input = schema(&fields); + let plan = CompiledPhysicalDag::from_operators( + [(0, InputContract::bounded(input.clone()))].into(), + Default::default(), + vec![0], + ) + .unwrap(); + let (result, _) = execute_batches(&plan, 1 << 20, 1, 42, false, |_, contract| { + Ok(BoundInput::Rows(Batch::try_new( + contract.schema.clone(), + vec![row.clone()], + )?)) + }) + .unwrap(); + let crate::query_engines::query_result::QueryResult::Vector(result) = result else { + panic!("vector expected") + }; + assert_eq!(result.values[0].value, expected); + } + } + // Count-like values are bound as integers only when the protocol sample is exact. #[test] fn integer_input_binding_preserves_type_and_rejects_rounding() { @@ -867,17 +904,25 @@ fn execute_batches( .fields .iter() .position(|field| field.name == SERIES_IDENTITY_COLUMN); - let value = batch - .schema() - .fields - .iter() - .position(|field| { - matches!( - field.dtype, - SummaryFamilyType::Plain(DataType::Float64 | DataType::Int64) - ) - }) - .ok_or_else(|| miss("physical output loses sample value"))?; + let fields = &batch.schema().fields; + let typed = |dtype: DataType| { + fields + .iter() + .enumerate() + .filter(|(_, field)| field.dtype == SummaryFamilyType::Plain(dtype.clone())) + .map(|(i, _)| i) + .collect::>() + }; + // The sample is the Float64 column; an Int64 count is the sample + // only when it is the sole numeric column. + let value = match ( + typed(DataType::Float64).as_slice(), + typed(DataType::Int64).as_slice(), + ) { + ([value, ..], _) => *value, + ([], [count]) => *count, + _ => return Err(miss("physical output loses sample value")), + }; for row in batch.rows() { let sample = &match row[value] { Value::Float64(sample) => sample, diff --git a/data_plane/src/storage_engines/sketch_db/current_series.rs b/data_plane/src/storage_engines/sketch_db/current_series.rs index e5d415023..f57d076c8 100644 --- a/data_plane/src/storage_engines/sketch_db/current_series.rs +++ b/data_plane/src/storage_engines/sketch_db/current_series.rs @@ -156,7 +156,7 @@ impl Population { .map(|(k, v)| (k.clone(), v.clone())) .collect(); self.remove(&labels); - // Bound the retained keys, tree nodes, and shared readout caches together. + // Bound the retained keys and tree nodes together. let bytes = 1024 + labels .iter() diff --git a/docs/developer_docs/query-engine/current-series-aggregations.md b/docs/developer_docs/query-engine/current-series-aggregations.md index 1c1eced54..a7119d50e 100644 --- a/docs/developer_docs/query-engine/current-series-aggregations.md +++ b/docs/developer_docs/query-engine/current-series-aggregations.md @@ -1,4 +1,4 @@ -# Shared current-series quantiles and TopK +# Shared current-series aggregations Backend-local PromQL workload compilation can export a maintained current-value alternative for `quantile(q, metric)`, `topk(k, metric)`, `sum`, `count` and @@ -50,6 +50,9 @@ A native exact alternative remains available for cost selection and execution fa unobserved generation or exceeded resource bounds cause native fallback. - State is in memory. Restart and generation replacement require warmup again; historical range queries continue to use native execution. +- Sum and average use Planner's aggregate operator, which adds values in order + without Prometheus' compensated summation; results can differ in the last + bits, and an average whose sum overflows is infinite. - Populations divide the configured retained-summary memory budget and cap series cardinality. Bounds include conservative space for labels and trees. The existing Remote Write adapter accepts finite sample values and stale markers.