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..5a623f1d6 100644 --- a/control_plane/src/physical/workload_cost.rs +++ b/control_plane/src/physical/workload_cost.rs @@ -834,14 +834,30 @@ 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) - .filter(|node| super::maintained_population::supported_node(node)) + 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| { + // 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) { @@ -1037,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() { @@ -1050,20 +1097,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..6448b4002 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,34 @@ 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 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 + && labels.is_some_and(|labels| { + labels == population.grouping.labels.iter().collect::>() + }); + 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..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,15 +904,31 @@ fn execute_batches( .fields .iter() .position(|field| field.name == SERIES_IDENTITY_COLUMN); - let value = batch - .schema() - .fields - .iter() - .position(|field| field.dtype == SummaryFamilyType::Plain(DataType::Float64)) - .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 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..f57d076c8 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); } @@ -160,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() @@ -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..a7119d50e 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 +# Shared current-series aggregations 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 @@ -47,13 +50,16 @@ 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, 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