diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index eb79859b..531a15a7 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -2139,7 +2139,7 @@ impl PhysicalCompiler { use std::hash::{Hash, Hasher}; let mut hash = std::collections::hash_map::DefaultHasher::new(); stable_workload_plan_id(&plan_materializations, &request.queries).hash(&mut hash); - "typed-local-residual-v1".hash(&mut hash); + "typed-local-residual-v2-range-max-index".hash(&mut hash); for query in &request.queries { format!("{:?}", query.post_asap).hash(&mut hash); } diff --git a/control_plane/src/physical/workload_cost.rs b/control_plane/src/physical/workload_cost.rs index 5938196d..a26499b3 100644 --- a/control_plane/src/physical/workload_cost.rs +++ b/control_plane/src/physical/workload_cost.rs @@ -207,13 +207,37 @@ pub fn manifest( // when no precomputed summary is installed. Deduplicate by source. for node in entry.nodes.values() { if let crate::query_plan::QueryPlanNode::Logical { - operator: crate::query_plan::logical::LogicalOperator::Scan { metric, .. }, + operator: + operator @ (crate::query_plan::logical::LogicalOperator::Scan { .. } + | crate::query_plan::logical::LogicalOperator::ReadRangeMaxIndex { + .. + }), .. } = node { - let metric = metric - .as_ref() - .ok_or_else(|| invalid("local raw scan requires named-source pricing"))?; + let metric = match operator { + crate::query_plan::logical::LogicalOperator::Scan { metric, .. } => metric + .as_ref() + .ok_or_else(|| invalid("local raw scan requires named-source pricing"))?, + crate::query_plan::logical::LogicalOperator::ReadRangeMaxIndex { + metric, + .. + } => metric, + _ => unreachable!(), + }; + if matches!( + operator, + crate::query_plan::logical::LogicalOperator::ReadRangeMaxIndex { .. } + ) { + for operation in ["build", "update", "residency", "retire"] { + add( + format!("range-max-index:{metric}:{operation}"), + json!({"operation": operation, "metric": metric, "index": "exact_per_series_range_max_v1"}), + "horizon", + 1.0, + ); + } + } let source = json!({"source": planner_types::pre_asap::Source::TimeSeries { metric: metric.clone() }, "location": "backend", "ingest": plan.precompute_plan.ingest}); add(format!("source:{}", source), source.clone(), "horizon", 1.0); for operation in ["build", "update", "residency", "retire"] { @@ -495,6 +519,52 @@ mod tests { snapshot } + #[test] + fn range_max_index_costs_share_state_across_filters_and_charge_retained_input() { + let mut snapshot = fixture(); + let entries = snapshot.query_workload.repeating_queries.as_mut().unwrap(); + entries[0].query = planner_types::workload::Query( + "max_over_time(service_retry_queue_depth{job=~\".+\"}[6h])".into(), + ); + entries[0].requirements.accuracy = planner_types::workload::AccuracyRequirement::Explicit( + planner_types::types::AccuracyTarget::Exact, + ); + let mut second = entries[0].clone(); + second.query = planner_types::workload::Query( + "max_over_time(service_retry_queue_depth{job=\"order-service\"}[6h])".into(), + ); + entries.push(second); + let (request, env) = snapshot.planning_request().unwrap(); + let plan = PhysicalCompiler.compile(request.clone(), env).unwrap(); + let costs = manifest(&plan, &request.queries).unwrap(); + assert_eq!( + costs + .components + .keys() + .filter(|id| id.starts_with("range-max-index:")) + .count(), + 4 + ); + assert_eq!( + costs + .components + .keys() + .filter(|id| id.starts_with("raw-state:")) + .count(), + 4 + ); + for operation in ["build", "update", "residency", "retire"] { + assert_eq!( + costs.components[&format!("range-max-index:service_retry_queue_depth:{operation}")] + .multiplicity, + 1.0 + ); + } + assert_eq!(plan.query_plan.entries.values().flat_map(|entry| entry.nodes.values()).filter(|node| + matches!(node, crate::query_plan::QueryPlanNode::Logical { operator: + crate::query_plan::logical::LogicalOperator::ReadRangeMaxIndex { .. }, .. })).count(), 2); + } + fn quoted() -> ( Vec, DeploymentEnvironment, diff --git a/control_plane/src/query_plan.rs b/control_plane/src/query_plan.rs index 1739e2c2..2d31a8b2 100644 --- a/control_plane/src/query_plan.rs +++ b/control_plane/src/query_plan.rs @@ -429,6 +429,18 @@ where let id = QueryNodeId(self.next_id); self.next_id += 1; self.seen.insert(identity, id); + if let Some(original) = &self.logical_source { + if let Some(operator) = logical::selected_range_max_index(original, node)? { + self.nodes.insert( + id, + QueryPlanNode::Logical { + operator, + inputs: vec![], + }, + ); + return Ok(id); + } + } let residual = match (&self.logical_source, &node.expr) { (Some(original), SummaryExpr::KeepPreAsap(expr)) => { Some(logical::residual_nodes(original, expr)?) diff --git a/control_plane/src/query_plan/logical.rs b/control_plane/src/query_plan/logical.rs index c22e0829..1aa883aa 100644 --- a/control_plane/src/query_plan/logical.rs +++ b/control_plane/src/query_plan/logical.rs @@ -12,6 +12,13 @@ use std::collections::BTreeMap; #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] #[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)] pub enum LogicalOperator { + /// Planner-selected exact per-series MinMax state, queried as max over an + /// exact event-time interval. This is an installed index, not a raw scan. + ReadRangeMaxIndex { + metric: String, + matchers: Vec, + range_ms: u64, + }, Scan { metric: Option, matchers: Vec, @@ -117,7 +124,7 @@ fn offset(value: &Option) -> Result { impl LogicalOperator { pub fn validate(&self, inputs: usize) -> Result<(), QueryPlanError> { let expected = match self { - Self::Scan { .. } => 0, + Self::Scan { .. } | Self::ReadRangeMaxIndex { .. } => 0, Self::Binary { .. } | Self::HistogramQuantile => 2, _ => 1, }; @@ -133,6 +140,14 @@ impl LogicalOperator { ) { return Err(invalid("zero range")); } + if let Self::ReadRangeMaxIndex { + metric, range_ms, .. + } = self + { + if metric.is_empty() || *range_ms == 0 || *range_ms > i64::MAX as u64 { + return Err(invalid("invalid exact range-max index contract")); + } + } if let Self::Subquery { range_ms, step_ms, .. } = self @@ -791,3 +806,113 @@ mod planner_workload_tests { } } } + +/// Preserve the exact original operator direction because MinMax family alone +/// does not distinguish min from max. The full Planner-node witness is required. +pub(super) fn selected_range_max_index( + original: &str, + node: &planner_types::post_asap::SummaryNode, +) -> Result, QueryPlanError> { + use planner_types::post_asap::{ExactKind, SummaryExpr, SummaryFamilyType}; + if !matches!( + &node.expr, + SummaryExpr::SummaryAgg { + family: SummaryFamilyType::ExactAggregate(ExactKind::MinMax, _), + reduction: planner_types::pre_asap::Reduction::PerEntity, + .. + } + ) { + return Ok(None); + } + let (root, nodes) = selected_residual_nodes(original, node)?; + let Some(QueryPlanNode::Logical { + operator: + LogicalOperator::Temporal { + operation: TemporalOperation::Max, + }, + inputs, + }) = nodes.get(&root) + else { + return Ok(None); + }; + if inputs.len() != 1 || nodes.len() != 2 { + return Ok(None); + } + let Some(QueryPlanNode::Logical { + operator: + LogicalOperator::Scan { + metric: Some(metric), + matchers, + range_ms: Some(range_ms), + offset_ms: 0, + }, + .. + }) = nodes.get(&inputs[0]) + else { + return Ok(None); + }; + Ok(Some(LogicalOperator::ReadRangeMaxIndex { + metric: metric.clone(), + matchers: matchers.clone(), + range_ms: *range_ms, + })) +} + +#[cfg(test)] +mod range_max_index_tests { + use super::*; + #[test] + fn real_gauge_queries_have_planner_authorized_exact_indexes() { + for (query, metric, range_ms) in [ + ( + r#"max_over_time(service_cache_refresh_lag_seconds{job="user-service"}[12h])"#, + "service_cache_refresh_lag_seconds", + 43_200_000, + ), + ( + r#"max_over_time(service_retry_queue_depth{job=~".+"}[6h])"#, + "service_retry_queue_depth", + 21_600_000, + ), + ( + r#"max_over_time(service_retry_queue_depth{job="order-service"}[6h])"#, + "service_retry_queue_depth", + 21_600_000, + ), + ] { + let original = crate::query_parser::parse_query_expr_canonical( + query, + planner_types::types::AccuracyTarget::Exact, + ) + .unwrap(); + let selected = crate::planner_selection::select_summary_default(&original).unwrap(); + let index = selected_range_max_index(query, &selected).unwrap().unwrap(); + assert!( + matches!(index, LogicalOperator::ReadRangeMaxIndex { metric: ref actual, range_ms: actual_range, .. } if actual == metric && actual_range == range_ms) + ); + index.validate(0).unwrap(); + assert!(index.validate(1).is_err()); + } + } + #[test] + fn min_and_shifted_or_nested_windows_do_not_become_max_indexes() { + for query in [ + "min_over_time(m[1m])", + "max_over_time(m[1m] offset 1m)", + "max_over_time((m + m)[1m:1s])", + ] { + let original = crate::query_parser::parse_query_expr_canonical( + query, + planner_types::types::AccuracyTarget::Exact, + ) + .unwrap(); + let selected = crate::planner_selection::select_summary_default(&original).unwrap(); + assert!( + selected_range_max_index(query, &selected) + .unwrap() + .is_none(), + "{query}" + ); + } + } +} diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index 21c5b7b0..dde9f945 100644 --- a/data_plane/src/drivers/ingest/prometheus_remote_write.rs +++ b/data_plane/src/drivers/ingest/prometheus_remote_write.rs @@ -203,13 +203,32 @@ impl PrometheusRemoteWriteReceiver { .ingest .physical_plan_snapshot() .ok_or_else(|| "active plan disappeared during input drain".to_string())?; - self.inner + let prepared = self + .inner .raw_store .snapshot( plan.precompute_plan.envelope.plan_id, plan.precompute_plan.envelope.plan_version, ) .map_err(|error| error.to_string())?; + let metrics = plan + .query_plan + .entries + .values() + .flat_map(|entry| entry.nodes.values()) + .filter_map(|node| match node { + control_plane::query_plan::QueryPlanNode::Logical { + operator: + control_plane::query_plan::logical::LogicalOperator::ReadRangeMaxIndex { + metric, + .. + }, + .. + } => Some(metric.clone()), + _ => None, + }) + .collect(); + prepared.prepare_range_max_indexes(&metrics); } Ok(()) } @@ -341,6 +360,16 @@ impl PrometheusRemoteWriteReceiver { } else { all_metrics = true; } + } else if let control_plane::query_plan::QueryPlanNode::Logical { + operator: + control_plane::query_plan::logical::LogicalOperator::ReadRangeMaxIndex { + metric, + .. + }, + .. + } = node + { + raw_metrics.insert(metric.clone()); } } } @@ -609,11 +638,12 @@ mod tests { } fn physical_config(streaming: StreamingConfig) -> HotReloadStreamingConfig { - physical_config_with_raw(streaming, false) + physical_config_with_raw(streaming, false, false) } fn physical_config_with_raw( streaming: StreamingConfig, retain_raw: bool, + indexed: bool, ) -> HotReloadStreamingConfig { use control_plane::physical::compiler::{ FrameIdentityContract, IngestContract, IngestProtocol, PlanEnvelope, PrecomputePlan, @@ -671,11 +701,19 @@ mod tests { nodes: std::collections::BTreeMap::from([( id, QueryPlanNode::Logical { - operator: control_plane::query_plan::logical::LogicalOperator::Scan { - metric: Some("requests_total".into()), - matchers: vec![], - range_ms: None, - offset_ms: 0, + operator: if indexed { + control_plane::query_plan::logical::LogicalOperator::ReadRangeMaxIndex { + metric: "requests_total".into(), + matchers: vec![], + range_ms: 60_000, + } + } else { + control_plane::query_plan::logical::LogicalOperator::Scan { + metric: Some("requests_total".into()), + matchers: vec![], + range_ms: None, + offset_ms: 0, + } }, inputs: vec![], }, @@ -715,6 +753,12 @@ mod tests { } fn configured_receiver_with_raw( retain_raw: bool, + ) -> (PrometheusRemoteWriteReceiver, mpsc::Receiver) { + configured_receiver_with_index(retain_raw, false) + } + fn configured_receiver_with_index( + retain_raw: bool, + indexed: bool, ) -> (PrometheusRemoteWriteReceiver, mpsc::Receiver) { use asap_types::enums::WindowKind; use asap_types::{AggregationConfig, AggregationType, KeyByLabelNames}; @@ -743,7 +787,7 @@ mod tests { router: SeriesRouter::new(vec![sender]), samples_ingested: AtomicU64::new(0), samples_blocked_by_schema_barrier: AtomicU64::new(0), - hot_reload_config: physical_config_with_raw(streaming, retain_raw), + hot_reload_config: physical_config_with_raw(streaming, retain_raw, indexed), pass_raw_samples: false, sketch_snapshots: dashmap::DashMap::new(), series_resolver: Arc::new(super::super::SeriesIdResolver::new()), @@ -812,6 +856,33 @@ mod tests { assert_eq!(native.raw_store().sample_count(), 0); } + #[tokio::test] + async fn finite_drain_eagerly_builds_installed_range_max_index() { + let (receiver, mut worker) = configured_receiver_with_index(true, true); + receiver.accept(&one_sample(3.0)).unwrap(); + assert_eq!( + receiver.raw_store().sample_count(), + 1, + "index-only plan must retain admitted input" + ); + assert_eq!(receiver.raw_store().range_max_index_bytes(), 0); + let handle = receiver.clone(); + let drain = tokio::spawn(async move { handle.drain().await }); + assert!(matches!( + worker.recv().await.unwrap(), + WorkerMessage::GroupSamples { .. } + )); + let WorkerMessage::Drain(reply) = worker.recv().await.unwrap() else { + panic!() + }; + reply.send(Ok(())).unwrap(); + drain.await.unwrap().unwrap(); + assert!( + receiver.raw_store().range_max_index_bytes() > 0, + "complete must include index construction before query calibration" + ); + } + // Closing finite input prevents writes racing behind the completion barrier. #[tokio::test] async fn finite_input_drain_seals_receiver_and_propagates_worker_failure() { diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index 76cddd3b..48a4ee73 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -5789,6 +5789,7 @@ async fn handle_store_metrics(State(state): State) -> axum::response:: "status": "success", "sid_count": timestamps.len(), "approx_resident_bytes": state.sketch_index.approx_resident_bytes(), + "exact_range_max_index_bytes": state.remote_write.as_ref().map(|receiver| receiver.raw_store().range_max_index_bytes()).unwrap_or(0), "raw_store_estimated_bytes": state.remote_write.as_ref().map(|receiver| receiver.raw_store().estimated_bytes()).unwrap_or(0), "raw_store_samples": state.remote_write.as_ref().map(|receiver| receiver.raw_store().sample_count()).unwrap_or(0), "earliest_timestamps_per_sid": timestamps}); diff --git a/data_plane/src/query_engines/asap_query_engine/logical_dag.rs b/data_plane/src/query_engines/asap_query_engine/logical_dag.rs index 476956dd..cf4f082f 100644 --- a/data_plane/src/query_engines/asap_query_engine/logical_dag.rs +++ b/data_plane/src/query_engines/asap_query_engine/logical_dag.rs @@ -26,6 +26,7 @@ struct Point { value: Option, } struct Series { + max_index: std::sync::OnceLock, labels: Labels, points: Vec, } @@ -85,15 +86,45 @@ impl PreparedSamples { } } points.dedup_by_key(|point| point.timestamp_ms); - series.push(Series { labels, points }); + series.push(Series { + labels, + points, + max_index: Default::default(), + }); } Ok(Self { series }) } + pub fn prepare_range_max_indexes(&self, metrics: &BTreeSet) { + for series in &self.series { + if series + .labels + .get("__name__") + .is_some_and(|metric| metrics.contains(metric)) + { + series.max_index.get_or_init(|| { + super::range_max_index::RangeMaxIndex::new( + series.points.iter().map(|point| point.value), + ) + }); + } + } + } + pub fn range_max_index_bytes(&self) -> usize { + self.series + .iter() + .filter_map(|series| series.max_index.get()) + .map(|index| index.estimated_bytes()) + .sum() + } pub fn estimated_bytes(&self) -> usize { self.series .iter() .map(|series| { std::mem::size_of::() + + series + .max_index + .get() + .map_or(0, |index| index.estimated_bytes()) + series.points.capacity() * std::mem::size_of::() + series .labels @@ -276,6 +307,11 @@ impl Result> Evaluator<' .ok_or_else(|| miss("missing logical input")) }; match operator { + LogicalOperator::ReadRangeMaxIndex { + metric, + matchers, + range_ms, + } => self.range_max_index(&metric, &matchers, range_ms, at), LogicalOperator::Scan { metric, matchers, @@ -419,6 +455,57 @@ impl Result> Evaluator<' } } } + fn range_max_index( + &mut self, + metric: &str, + matchers: &[LabelMatcher], + range_ms: u64, + at: i64, + ) -> Result { + for matcher in matchers { + if matches!(matcher.operation, LabelMatch::Regex | LabelMatch::NotRegex) + && !self.regexes.contains_key(&matcher.value) + { + let regex = regex::Regex::new(&format!("(?s)^(?:{})$", matcher.value)) + .map_err(|e| miss(format!("unsupported regex: {e}")))?; + self.regexes.insert(matcher.value.clone(), regex); + } + } + let start = at + .checked_sub(i64::try_from(range_ms).map_err(|_| miss("range overflow"))?) + .ok_or_else(|| miss("range overflow"))?; + let mut values = Vec::new(); + for series in self.series { + if series.labels.get("__name__").map(String::as_str) != Some(metric) + || !matchers.iter().all(|matcher| { + let value = series + .labels + .get(&matcher.name) + .map(String::as_str) + .unwrap_or(""); + match matcher.operation { + LabelMatch::Equal => value == matcher.value, + LabelMatch::NotEqual => value != matcher.value, + LabelMatch::Regex => self.regexes[&matcher.value].is_match(value), + LabelMatch::NotRegex => !self.regexes[&matcher.value].is_match(value), + } + }) + { + continue; + } + let lo = series.points.partition_point(|p| p.timestamp_ms <= start); + let hi = series.points.partition_point(|p| p.timestamp_ms <= at); + let index = series.max_index.get_or_init(|| { + super::range_max_index::RangeMaxIndex::new(series.points.iter().map(|p| p.value)) + }); + // Count only real index access, never a generic scan or successful binding. + self.stats.summary_readout_evaluations += 1; + if let Some(value) = index.query(lo, hi) { + values.push((no_name(series.labels.clone()), value)); + } + } + Ok(Value::Vector(values)) + } fn scan( &mut self, metric: Option<&str>, @@ -763,6 +850,180 @@ mod tests { }; result.values.into_iter().map(|p| p.value).collect() } + #[test] + fn indexed_max_preserves_filters_labels_boundaries_and_real_readout_provenance() { + let mut samples = vec![ + sample( + "service_retry_queue_depth", + "user-service", + 1_000, + Some(999.), + ), + sample("service_retry_queue_depth", "user-service", 1_001, Some(4.)), + sample("service_retry_queue_depth", "user-service", 2_056, Some(8.)), + sample( + "service_retry_queue_depth", + "order-service", + 2_000, + Some(30.), + ), + ]; + let mut extra = sample( + "service_retry_queue_depth", + "user-service", + 1_500, + Some(11.), + ); + extra.labels.insert("instance".into(), "other".into()); + samples.push(extra); + let prepared = PreparedSamples::new(&samples).unwrap(); + assert_eq!(prepared.range_max_index_bytes(), 0); + for selector in [ + r#"job="user-service""#, + r#"job=~".+""#, + r#"job!="order-service",absent="""#, + r#"job!~"order.*""#, + ] { + let query = format!("max_over_time(service_retry_queue_depth{{{selector}}}[1056ms])"); + let raw = entry(&query); + let mut indexed = raw.clone(); + let scan = raw + .nodes + .values() + .find_map(|node| match node { + QueryPlanNode::Logical { + operator: + LogicalOperator::Scan { + metric: Some(metric), + matchers, + range_ms: Some(range_ms), + .. + }, + .. + } => Some(LogicalOperator::ReadRangeMaxIndex { + metric: metric.clone(), + matchers: matchers.clone(), + range_ms: *range_ms, + }), + _ => None, + }) + .unwrap(); + indexed.nodes = [( + indexed.root, + QueryPlanNode::Logical { + operator: scan, + inputs: vec![], + }, + )] + .into(); + for at in [2056, 2057, 2500, 4000] { + let (expected, _) = + execute_prepared_with_stats(&raw, &prepared, at, |_, _| unreachable!()) + .unwrap(); + let (actual, stats) = + execute_prepared_with_stats(&indexed, &prepared, at, |_, _| unreachable!()) + .unwrap(); + assert_eq!( + serde_json::to_value(actual).unwrap(), + serde_json::to_value(expected).unwrap(), + "{query} at {at}" + ); + assert_eq!(stats.raw_scan_evaluations, 0); + assert!(stats.summary_readout_evaluations > 0); + } + } + let bytes = prepared.range_max_index_bytes(); + assert!(bytes > 0 && bytes < prepared.estimated_bytes()); + } + + #[test] + fn real_planner_max_index_matches_exact_at_56_seconds_and_rebuilds_after_input() { + use crate::query_engines::raw_store::RawSampleStore; + use control_plane::physical::compiler::BackendLocalPlanningSnapshot; + let query = r#"max_over_time(service_cache_refresh_lag_seconds{job="user-service"}[12h])"#; + let mut snapshot: BackendLocalPlanningSnapshot = serde_json::from_str(include_str!( + "../../../../docs/examples/asapquery-planning-snapshot.json" + )) + .unwrap(); + let q = &mut snapshot.query_workload.repeating_queries.as_mut().unwrap()[0]; + q.query = planner_types::workload::Query(query.into()); + q.requirements.accuracy = planner_types::workload::AccuracyRequirement::Explicit( + planner_types::types::AccuracyTarget::Exact, + ); + // Compile the evidence-collection candidate; v1 unquoted startup keeps + // its native compatibility policy, while production v2 prices candidates. + let (candidate, environment) = snapshot.planning_request().unwrap(); + let plan = control_plane::physical::compiler::PhysicalCompiler + .compile(candidate, environment) + .unwrap(); + let installed = plan.query_plan.lookup(query).unwrap(); + assert!(installed.nodes.values().any(|n| matches!( + n, + QueryPlanNode::Logical { + operator: LogicalOperator::ReadRangeMaxIndex { .. }, + .. + } + ))); + let at = 1_788_891_296_000_i64; + let start = at - 43_200_000; + let metric = "service_cache_refresh_lag_seconds"; + let samples = [ + sample(metric, "user-service", start, Some(999.)), + sample(metric, "user-service", start + 1, Some(4.)), + sample(metric, "user-service", at, Some(8.)), + sample(metric, "order-service", at, Some(500.)), + ]; + let metrics = BTreeSet::from([metric.into()]); + let store = RawSampleStore::default(); + store.append_admitted(1, 1, &samples, &metrics, false); + let first = store.snapshot(1, 1).unwrap(); + first.prepare_range_max_indexes(&metrics); + let bytes = store.range_max_index_bytes(); + assert!(bytes > 0); + let execute_index = |data: &PreparedSamples| { + execute_prepared_with_stats(installed, data, at as u64, |_, _| unreachable!()).unwrap() + }; + let (actual, stats) = execute_index(&first); + let expected = execute(&entry(query), &samples, at as u64).unwrap(); + assert_eq!( + serde_json::to_value(actual).unwrap(), + serde_json::to_value(expected).unwrap() + ); + assert_eq!(stats.raw_scan_evaluations, 0); + assert_eq!(stats.summary_readout_evaluations, 1); + assert_eq!( + store.range_max_index_bytes(), + bytes, + "query must reuse eager state" + ); + store.append_admitted( + 1, + 1, + &[sample(metric, "user-service", at - 1, Some(20.))], + &metrics, + false, + ); + assert_eq!( + store.range_max_index_bytes(), + 0, + "admission must invalidate prior index" + ); + let next = store.snapshot(1, 1).unwrap(); + assert!(!std::sync::Arc::ptr_eq(&first, &next)); + next.prepare_range_max_indexes(&metrics); + let (QueryResult::Vector(updated), _) = execute_index(&next) else { + panic!() + }; + assert_eq!(updated.values[0].value, 20.); + let (QueryResult::Vector(old), _) = execute_index(&first) else { + panic!() + }; + assert_eq!( + old.values[0].value, 8., + "in-flight snapshot remains coherent" + ); + } + // Regex alternation remains fully anchored and missing labels compare as empty. #[test] fn anchored_matchers_and_missing_labels() { diff --git a/data_plane/src/query_engines/asap_query_engine/mod.rs b/data_plane/src/query_engines/asap_query_engine/mod.rs index 963d74bc..1bba0f66 100644 --- a/data_plane/src/query_engines/asap_query_engine/mod.rs +++ b/data_plane/src/query_engines/asap_query_engine/mod.rs @@ -26,3 +26,5 @@ pub use crate::storage_engines::sketch_db::query as asap_tier; pub mod tests; pub use engine::ASAPQueryEngine; + +mod range_max_index; diff --git a/data_plane/src/query_engines/asap_query_engine/range_max_index.rs b/data_plane/src/query_engines/asap_query_engine/range_max_index.rs new file mode 100644 index 00000000..fabaae80 --- /dev/null +++ b/data_plane/src/query_engines/asap_query_engine/range_max_index.rs @@ -0,0 +1,108 @@ +//! Exact range-max state for one complete source-label set. +//! +//! Leaves retain timestamp order. Query bounds are sample offsets computed with +//! (start, end] partition points; no pane rounding or counter semantics apply. +#[derive(Debug)] +pub(crate) struct RangeMaxIndex { + tree: Vec>, + width: usize, + len: usize, +} +fn merge(left: Option, right: Option) -> Option { + match (left, right) { + (None, value) | (value, None) => value, + (Some(a), Some(b)) => Some(if a.is_nan() || b > a { b } else { a }), + } +} +impl RangeMaxIndex { + pub(crate) fn new(values: impl ExactSizeIterator>) -> Self { + let len = values.len(); + let width = len.max(1).next_power_of_two(); + let mut tree = vec![None; width * 2]; + for (i, value) in values.enumerate() { + tree[width + i] = value.filter(|v| v.to_bits() != 0x7ff0000000000002); + } + for i in (1..width).rev() { + tree[i] = merge(tree[i * 2], tree[i * 2 + 1]); + } + Self { tree, width, len } + } + pub(crate) fn query(&self, start: usize, end: usize) -> Option { + assert!(start <= end && end <= self.len); + let (mut lo, mut hi) = (start + self.width, end + self.width); + let (mut left, mut right) = (None, None); + while lo < hi { + if lo % 2 == 1 { + left = merge(left, self.tree[lo]); + lo += 1; + } + if hi % 2 == 1 { + hi -= 1; + right = merge(self.tree[hi], right); + } + lo /= 2; + hi /= 2; + } + merge(left, right) + } + pub(crate) fn estimated_bytes(&self) -> usize { + self.tree.capacity() * std::mem::size_of::>() + } +} +#[cfg(test)] +mod tests { + use super::*; + #[test] + fn every_partial_interval_matches_ordered_prometheus_max_semantics() { + let values = [ + None, + Some(-0.0), + Some(0.0), + Some(f64::NAN), + Some(-4.0), + Some(7.0), + Some(f64::from_bits(0x7ff0000000000002)), + Some(f64::NEG_INFINITY), + ]; + let index = RangeMaxIndex::new(values.into_iter()); + for start in 0..=values.len() { + for end in start..=values.len() { + let expected = values[start..end] + .iter() + .copied() + .map(|v| v.filter(|v| v.to_bits() != 0x7ff0000000000002)) + .fold(None, merge); + assert_eq!( + index.query(start, end).map(f64::to_bits), + expected.map(f64::to_bits), + "{start}..{end}" + ); + } + } + assert!(index.estimated_bytes() > 0); + } + #[test] + fn signed_zero_and_nan_follow_first_non_nan_max_in_timestamp_order() { + let index = RangeMaxIndex::new( + [ + Some(f64::NAN), + Some(-0.0), + Some(0.0), + Some(f64::NAN), + Some(-2.0), + ] + .into_iter(), + ); + assert!(index.query(0, 1).unwrap().is_nan()); + assert_eq!(index.query(0, 3).unwrap().to_bits(), (-0.0_f64).to_bits()); + assert_eq!(index.query(2, 4).unwrap().to_bits(), 0.0_f64.to_bits()); + assert_eq!(index.query(3, 5), Some(-2.0)); + } + #[test] + fn empty_and_all_nan_ranges_are_distinct() { + let empty = RangeMaxIndex::new([].into_iter()); + assert_eq!(empty.query(0, 0), None); + let nan = RangeMaxIndex::new([Some(f64::NAN), Some(f64::NAN)].into_iter()); + assert!(nan.query(0, 2).unwrap().is_nan()); + } +} diff --git a/data_plane/src/query_engines/raw_store.rs b/data_plane/src/query_engines/raw_store.rs index 99890bf5..194c9efe 100644 --- a/data_plane/src/query_engines/raw_store.rs +++ b/data_plane/src/query_engines/raw_store.rs @@ -95,6 +95,14 @@ impl RawSampleStore { Ok(prepared) } + pub fn range_max_index_bytes(&self) -> usize { + let state = self.generation.lock().unwrap_or_else(|e| e.into_inner()); + state + .prepared + .as_ref() + .map_or(0, |p| p.range_max_index_bytes()) + } + pub fn sample_count(&self) -> usize { self.generation .lock()