From 06da404bca6eba8fdc972a993346bf39c18bd2eb Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 12:36:12 +0000 Subject: [PATCH 1/3] feat(types): admit raw inputs beside stored native batches Replace the no-common-snapshot reject with a bounded-lag admission: query-time raw inputs may share a physical program with stored native batches, whose lag the executor checks. Add the stored window grid helpers used to find the newest complete window. Co-Authored-By: Claude Opus 5.5 --- crates/asap_types/src/query_plan.rs | 70 ++++++++++++++++++++++ crates/asap_types/src/query_plan/native.rs | 58 ++++++++++++------ 2 files changed, 111 insertions(+), 17 deletions(-) diff --git a/crates/asap_types/src/query_plan.rs b/crates/asap_types/src/query_plan.rs index 8aa371acb..9cf1c7cc5 100644 --- a/crates/asap_types/src/query_plan.rs +++ b/crates/asap_types/src/query_plan.rs @@ -756,6 +756,25 @@ impl MaterializationBinding { } } } + + /// Interval between the ends of successive complete stored windows. + pub fn slide_ms(&self) -> u64 { + self.full_window_slide_ms.unwrap_or(self.window_ms) + } + + /// End of the newest window on this output's grid that ends at or before + /// `at_ms`; None when the grid is unknown or that window starts before 0. + pub fn latest_window_end_at_or_before(&self, at_ms: u64) -> Option { + let origin = i128::from(self.pane_origin_ms?); + let (window, slide) = (i128::from(self.window_ms), i128::from(self.slide_ms())); + if window == 0 || slide == 0 { + return None; + } + // Pane and full-window ends both lie on origin + window + n * slide. + let at = i128::from(at_ms); + let end = at - (at - origin - window).rem_euclid(slide); + u64::try_from(end).ok().filter(|end| *end >= self.window_ms) + } } #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] @@ -1096,3 +1115,54 @@ mod retired_plan_tests { } } } + +#[cfg(test)] +mod window_grid_tests { + use super::*; + + fn binding(window_ms: u64, slide: Option, origin: Option) -> MaterializationBinding { + MaterializationBinding { + stored_output_reference: serde_json::from_value(serde_json::json!({ + "stored_output_id": 1, + "definition_id": format!("sds-v1:{}", "a".repeat(64)), + })) + .unwrap(), + full_window_slide_ms: slide, + materialization: StoredOutputId(1), + output_grouping: PhysicalGrouping::Reduce(vec![]), + item_labels: vec![], + window_ms, + pane_origin_ms: origin, + readout_lookback_ms: Some(window_ms), + } + } + + // The newest complete window end follows the pane or full-window grid at any + // origin, and every returned window is one the grid can store. + #[test] + fn latest_window_end_follows_the_stored_grid() { + let panes = binding(60_000, None, Some(0)); + assert_eq!(panes.slide_ms(), 60_000); + assert_eq!(panes.latest_window_end_at_or_before(120_000), Some(120_000)); + assert_eq!(panes.latest_window_end_at_or_before(179_999), Some(120_000)); + assert_eq!(panes.latest_window_end_at_or_before(59_999), None); + let sliding = binding(300_000, Some(60_000), Some(5_000)); + assert_eq!(sliding.slide_ms(), 60_000); + let end = sliding.latest_window_end_at_or_before(1_000_000).unwrap(); + assert_eq!(end, 965_000); + assert!(sliding.covers_range(end - 300_000, end)); + assert_eq!( + binding(60_000, None, None).latest_window_end_at_or_before(120_000), + None + ); + let shifted = binding(60_000, None, Some(-30_000)); + assert_eq!( + shifted.latest_window_end_at_or_before(100_000), + Some(90_000) + ); + assert_eq!( + binding(60_000, None, Some(500_000)).latest_window_end_at_or_before(100_000), + Some(80_000) + ); + } +} diff --git a/crates/asap_types/src/query_plan/native.rs b/crates/asap_types/src/query_plan/native.rs index 6448b4002..df10b6847 100644 --- a/crates/asap_types/src/query_plan/native.rs +++ b/crates/asap_types/src/query_plan/native.rs @@ -116,6 +116,32 @@ impl QueryPlanEntry { } } + /// A query-time raw input: a named range selector read from the raw-series + /// endpoint at the evaluation time. + fn is_query_time_raw(&self, input: &QueryNodeId) -> bool { + matches!( + self.nodes.get(input), + Some(QueryPlanNode::Logical { + operator: query_time::QueryTimeOperator::Scan { + metric: Some(_), + range_ms: Some(_), + .. + }, + .. + }) + ) + } + + /// True when the physical inputs include both query-time raw inputs and + /// stored inputs. + pub fn mixes_raw_and_stored_inputs(&self) -> bool { + self.physical_vector_binding() + .is_some_and(|(inputs, _, _)| { + inputs.iter().any(|input| self.is_query_time_raw(input)) + && !inputs.iter().all(|input| self.is_query_time_raw(input)) + }) + } + /// Whether the root physical result drops `__name__` from series identities. pub fn drops_metric_name(&self) -> bool { matches!( @@ -142,24 +168,22 @@ impl QueryPlanEntry { { return Err(invalid("invalid physical vector source mapping or budget")); } - let raw = |input: &QueryNodeId| { - matches!( - self.nodes.get(input), - Some(QueryPlanNode::Logical { - operator: query_time::QueryTimeOperator::Scan { - metric: Some(_), - range_ms: Some(_), - .. - }, - .. - }) - ) - }; - // Raw samples are read from the external endpoint at query time; like - // exact cuts, they share no snapshot with installed summary state. - if inputs.iter().any(raw) && !inputs.iter().all(raw) { + let raw = |input: &QueryNodeId| self.is_query_time_raw(input); + // Only a stored native batch has a lag the executor can check against + // the evaluation time of the raw inputs. Its admitted families (Sum and + // heap sketches, checked below) read out independently of the + // evaluation time, so a lagged batch is read as of its own window. + if inputs.iter().any(raw) + && inputs.iter().any(|input| { + !raw(input) + && !matches!( + self.nodes.get(input), + Some(QueryPlanNode::ReadMaterialization { .. }) + ) + }) + { return Err(invalid( - "query-time raw inputs cannot be mixed with installed state", + "query-time raw inputs mix only with stored native batches", )); } for input in inputs { From 23a963a2e64516d580fac7fd293100a3266d31b6 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 12:36:13 +0000 Subject: [PATCH 2/3] feat: bind mixed raw and stored inputs under a bounded lag Raw inputs are read at t_q; each stored input uses its newest complete window t_s and is admitted only when t_q - t_s <= max_lag (default one slide of that output, override --max-stored-input-lag-ms). Beyond the bound the query misses to the exact fallback. The observed lag is traced and reported in execution stats. Co-Authored-By: Claude Opus 5.5 --- data_plane/src/main.rs | 7 + .../query_engines/asap_query_engine/engine.rs | 66 ++- .../asap_query_engine/logical_dag.rs | 3 + .../logical_dag/native_values.rs | 145 ++++-- .../asap_query_engine/raw_source.rs | 440 ++++++++++++++++++ 5 files changed, 621 insertions(+), 40 deletions(-) diff --git a/data_plane/src/main.rs b/data_plane/src/main.rs index d5b9e5d7a..4be5732d2 100644 --- a/data_plane/src/main.rs +++ b/data_plane/src/main.rs @@ -159,6 +159,12 @@ struct Args { #[arg(long)] disable_query_forwarding: bool, + /// Largest staleness of stored state combined with raw series read at the + /// query time, applied to every stored output; beyond it the query uses + /// the exact fallback. Unset: one slide interval of each stored output. + #[arg(long)] + max_stored_input_lag_ms: Option, + /// Database path (currently unused, kept for compatibility) #[arg(long, default_value = "sketchdb.db")] db_path: String, @@ -762,6 +768,7 @@ async fn main() -> Result<()> { .with_sketch_index(summary_store.clone()) .with_active_physical_plan(active_physical_plan.clone()) .with_query_forwarding_policy(query_forwarding_policy) + .with_max_stored_input_lag_ms(args.max_stored_input_lag_ms) .with_exact_subquery_endpoint(args.prometheus_server.clone()) .with_metricsql_exact_subquery_endpoint(args.victoriametrics_url.clone()); 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 46f798d5f..292fa4ca7 100644 --- a/data_plane/src/query_engines/asap_query_engine/engine.rs +++ b/data_plane/src/query_engines/asap_query_engine/engine.rs @@ -78,6 +78,36 @@ mod readiness_coverage_tests { } } +#[cfg(test)] +mod stored_input_lag_tests { + // Only an execution that bound stored state beside raw inputs reports its lag. + #[test] + fn logical_execution_reports_stored_input_lag_when_mixed() { + use crate::query_engines::query_result::QueryResult; + let lag_warnings = |lag| { + let mut result = QueryResult::vector(vec![], 0); + let stats = super::super::logical_dag::ExecutionStats { + stored_input_lag_ms: lag, + ..Default::default() + }; + super::annotate_logical_execution(&mut result, &stats); + let QueryResult::Vector(result) = result else { + unreachable!() + }; + result + .warnings + .into_iter() + .filter(|warning| warning.starts_with("asap_stored_input_lag_ms:")) + .collect::>() + }; + assert_eq!( + lag_warnings(Some(40_000)), + ["asap_stored_input_lag_ms:40000"] + ); + assert!(lag_warnings(None).is_empty()); + } +} + #[cfg(test)] mod forwarding_policy_tests { use super::ASAPQueryEngine; @@ -198,6 +228,9 @@ pub struct ASAPQueryEngine { query_forwarding_policy: crate::query_engines::QueryForwardingPolicy, exact_subquery_client: reqwest::Client, execution_limits: asap_physical_operators::dag::Limits, + /// Deployment override of the staleness bound for stored inputs bound + /// beside query-time raw inputs; None uses each stored output's slide. + max_stored_input_lag_ms: Option, } impl ASAPQueryEngine { @@ -278,6 +311,7 @@ impl ASAPQueryEngine { Self { prometheus_scrape_interval, execution_limits: Default::default(), + max_stored_input_lag_ms: None, summary_store: None, active_physical_plan: None, exact_subquery_endpoint: None, @@ -307,6 +341,10 @@ impl ASAPQueryEngine { self.query_forwarding_policy = policy; self } + pub fn with_max_stored_input_lag_ms(mut self, max_lag_ms: Option) -> Self { + self.max_stored_input_lag_ms = max_lag_ms; + self + } pub fn with_execution_limits(mut self, limits: asap_physical_operators::dag::Limits) -> Self { self.execution_limits = limits; self @@ -535,13 +573,31 @@ impl ASAPQueryEngine { .as_deref() .filter(|_| self.query_forwarding_policy.allows_external_queries()) .map(|endpoint| (&self.exact_subquery_client, endpoint)); - super::logical_dag::native_values::execute_stored( - entry, + let (plan_id, plan_version) = ( physical.query_plan.plan_id, physical.query_plan.plan_version, - index, + ); + super::logical_dag::native_values::execute_stored( + entry, raw_endpoint, + self.max_stored_input_lag_ms, at, + &mut |binding, window, schema, max_bytes| { + let store = index.ok_or("summary store unavailable")?; + let address = asap_types::sds::StoredSummaryKey { + plan_id, + plan_version, + stored_output_id: binding.stored_output_reference.stored_output_id, + population: std::collections::BTreeMap::new(), + window, + }; + store.read_bound_native_summary( + &address, + &binding.stored_output_reference, + schema, + max_bytes, + ) + }, ) } else { super::logical_dag::execute_installed(entry, leaves, at, |root, evaluation_ms| { @@ -737,6 +793,7 @@ impl ASAPQueryEngine { total.remote_evaluations += stats.remote_evaluations; total.remote_rpcs += stats.remote_rpcs; total.remote_branch_evaluations += stats.remote_branch_evaluations; + total.stored_input_lag_ms = total.stored_input_lag_ms.max(stats.stored_input_lag_ms); let _step_result = crate::query_engines::request::reserve(result.retained_bytes())?; let QueryResult::Vector(result) = result else { return Err(EngineError::capability_miss( @@ -1059,6 +1116,9 @@ fn annotate_logical_execution( .into(), ); } + if let Some(lag) = stats.stored_input_lag_ms { + warnings.push(format!("asap_stored_input_lag_ms:{lag}")); + } warnings.push(format!( "asap_logical_stats:raw={},summary={},memo_hits={},remote={},remote_rpcs={},remote_branches={}", stats.raw_scan_evaluations, 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 9fa65f67d..505424279 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 @@ -57,6 +57,9 @@ pub struct ExecutionStats { pub remote_evaluations: usize, pub remote_rpcs: usize, pub remote_branch_evaluations: usize, + /// Largest `t_q - t_s` over stored inputs bound beside query-time raw + /// inputs; None when the program has no such mix. + pub stored_input_lag_ms: Option, } fn miss(detail: impl Into) -> EngineError { EngineError::capability_miss("installed_logical_dag", detail) 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 4447f17b2..5b05d09e2 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 @@ -774,11 +774,19 @@ where pub(in crate::query_engines::asap_query_engine) fn execute_stored( entry: &asap_types::query_plan::QueryPlanEntry, - plan_id: u64, - plan_version: u64, - store: Option<&crate::storage_engines::sketch_db::index::SketchStore>, raw_endpoint: Option<(&reqwest::Client, &str)>, + max_stored_lag_ms: Option, at: u64, + // Reads one stored native batch for a window under its installed binding. + read: &mut dyn FnMut( + &asap_types::query_plan::MaterializationBinding, + asap_types::sds::HalfOpenTimeRange, + Schema, + usize, + ) -> Result< + Batch, + crate::storage_engines::sketch_db::index::NativeReadError, + >, ) -> Result< ( crate::query_engines::query_result::QueryResult, @@ -793,7 +801,9 @@ pub(in crate::query_engines::asap_query_engine) fn execute_stored( .recover_vector_physical_dag() .map_err(|e| miss(e.to_string()))?; let drop_metric_name = entry.drops_metric_name(); - execute_batches( + let mixed = entry.mixes_raw_and_stored_inputs(); + let mut observed_lag = None; + let (result, mut stats) = execute_batches( &program, max_bytes, inputs.len(), @@ -818,41 +828,102 @@ pub(in crate::query_engines::asap_query_engine) fn execute_stored( else { return Err(miss("native stored source has no deployed summary binding")); }; - let store = store.ok_or_else(|| { - EngineError::capability_miss("native_stored", "summary store unavailable") - })?; - let end = i64::try_from(at).map_err(|_| miss("native timestamp overflow"))?; - let start = at - .checked_sub(binding.window_ms) - .ok_or_else(|| miss("native window underflow"))?; - let address = asap_types::sds::StoredSummaryKey { - plan_id, - plan_version, - stored_output_id: binding.stored_output_reference.stored_output_id, - population: std::collections::BTreeMap::new(), - window: asap_types::sds::HalfOpenTimeRange { - start_ms: start as i64, - end_ms: end, - }, + let mut read = |end: u64| { + let start = end + .checked_sub(binding.window_ms) + .ok_or("native window underflow")?; + let window = asap_types::sds::HalfOpenTimeRange { + start_ms: i64::try_from(start).map_err(|_| "native timestamp overflow")?, + end_ms: i64::try_from(end).map_err(|_| "native timestamp overflow")?, + }; + read(binding, window, schema.clone(), max_bytes as usize) }; - store - .read_bound_native_summary( - &address, - &binding.stored_output_reference, - schema.clone(), - max_bytes as usize, - ) - .map(BoundInput::Rows) - .map_err(|error| match error { - crate::storage_engines::sketch_db::index::NativeReadError::Unavailable( - message, - ) => miss(message), - crate::storage_engines::sketch_db::index::NativeReadError::Physical(error) => { - EngineError::from(error) - } - }) + if !mixed { + return read(at).map(BoundInput::Rows).map_err(native_read_error); + } + let (batch, lag) = read_within_lag(binding, at, max_stored_lag_ms, &mut read)?; + observed_lag = observed_lag.max(Some(lag)); + Ok(BoundInput::Rows(batch)) }, - ) + )?; + stats.stored_input_lag_ms = observed_lag; + Ok((result, stats)) +} + +fn native_read_error( + error: crate::storage_engines::sketch_db::index::NativeReadError, +) -> EngineError { + match error { + crate::storage_engines::sketch_db::index::NativeReadError::Unavailable(message) => { + miss(message) + } + crate::storage_engines::sketch_db::index::NativeReadError::Physical(error) => { + EngineError::from(error) + } + } +} + +/// Bind a stored input beside raw inputs read at `at`: its newest complete +/// window, ending at `t_s`, is admitted only when `at - t_s` is within the +/// bound, which defaults to one slide of the stored output. Otherwise the +/// query misses and takes the exact fallback rather than mixing staler state. +/// Returns the batch and its observed lag. +fn read_within_lag( + binding: &asap_types::query_plan::MaterializationBinding, + at: u64, + max_lag_ms: Option, + read: &mut dyn FnMut( + u64, + ) -> Result< + Batch, + crate::storage_engines::sketch_db::index::NativeReadError, + >, +) -> Result<(Batch, u64), EngineError> { + use crate::storage_engines::sketch_db::index::NativeReadError; + let max_lag = max_lag_ms.unwrap_or_else(|| binding.slide_ms()); + let mut end = binding + .latest_window_end_at_or_before(at) + .ok_or_else(|| miss("stored input has no complete window grid before the query time"))?; + let mut unavailable = String::new(); + loop { + let lag = at - end; + if lag > max_lag { + tracing::debug!( + stored_output = ?binding.stored_output_reference.stored_output_id, + query_time_ms = at, + lag_ms = lag, + max_lag_ms = max_lag, + "stored input exceeds bounded lag; using exact fallback" + ); + return Err(miss(format!( + "stored input lag {lag}ms exceeds bound {max_lag}ms at {at}: {unavailable}" + ))); + } + match read(end) { + Ok(batch) => { + tracing::debug!( + stored_output = ?binding.stored_output_reference.stored_output_id, + query_time_ms = at, + stored_watermark_ms = end, + lag_ms = lag, + max_lag_ms = max_lag, + "stored input admitted within bounded lag" + ); + return Ok((batch, lag)); + } + // The window is not complete yet; an older one may still be in bound. + Err(NativeReadError::Unavailable(message)) => unavailable = message, + Err(error) => return Err(native_read_error(error)), + } + end = match end.checked_sub(binding.slide_ms()) { + Some(older) if older >= binding.window_ms => older, + _ => { + return Err(miss(format!( + "stored input has no complete window within its lag bound: {unavailable}" + ))) + } + }; + } } /// Stored and protocol inputs are read before execution; query-time raw diff --git a/data_plane/src/query_engines/asap_query_engine/raw_source.rs b/data_plane/src/query_engines/asap_query_engine/raw_source.rs index 5928918d2..3baf830e2 100644 --- a/data_plane/src/query_engines/asap_query_engine/raw_source.rs +++ b/data_plane/src/query_engines/asap_query_engine/raw_source.rs @@ -659,6 +659,13 @@ mod tests { ]) ); assert_eq!(seen.lock().unwrap().len(), 1); + assert!( + !result + .warnings + .iter() + .any(|warning| warning.starts_with("asap_stored_input_lag_ms:")), + "a raw-only program has no stored lag" + ); server.abort(); } @@ -728,4 +735,437 @@ mod tests { "{error}" ); } + + mod mixed_inputs { + //! Hand-built programs `sum_over_time(m[5m]) + `: + //! the raw input is read at the query time, the stored one at its + //! newest complete window within the lag bound. + use super::*; + use crate::storage_engines::sketch_db::index::NativeReadError; + use asap_physical_operators::{ + operators::{Operator, ReadoutQuery}, + physical_planner::promql_rows::encode_series_identity, + summary_kernels::exact::{ExactAccumulator, ExactReadout}, + }; + use asap_types::sds::HalfOpenTimeRange; + use planner_types::post_asap::{ExactKind, ExactParams, SummaryField, SummarySchema}; + + const WINDOW: u64 = 300_000; + const SLIDE: u64 = 60_000; + /// Newest stored window end at or before AT on the (origin 0) grid. + const LATEST: u64 = 960_000; + const AT_MS: u64 = AT as u64; + + fn field(name: &str, dtype: SummaryFamilyType) -> SummaryField { + SummaryField { + name: name.into(), + dtype, + nullable: false, + } + } + fn sum_family() -> SummaryFamilyType { + SummaryFamilyType::ExactAggregate(ExactKind::Sum, ExactParams::Sum) + } + fn raw_schema() -> Schema { + Arc::new(SummarySchema { + fields: vec![ + field( + SERIES_IDENTITY_COLUMN, + SummaryFamilyType::Plain(DataType::Utf8), + ), + field("time", SummaryFamilyType::Plain(DataType::Timestamp)), + field("value", SummaryFamilyType::Plain(DataType::Float64)), + ], + time_index: Some(1), + }) + } + fn stored_schema() -> Schema { + Arc::new(SummarySchema { + fields: vec![ + field( + SERIES_IDENTITY_COLUMN, + SummaryFamilyType::Plain(DataType::Utf8), + ), + field("state", sum_family()), + ], + time_index: None, + }) + } + fn stored_readout() -> Operator { + Operator::readout( + stored_schema(), + 1, + ReadoutQuery::Exact(ExactReadout { + statistic: asap_physical_operators::Statistic::Sum, + lookback_ms: None, + }), + ) + .unwrap() + } + fn binding(slide: u64) -> MaterializationBinding { + MaterializationBinding { + stored_output_reference: serde_json::from_value(serde_json::json!({ + "stored_output_id": 7, + "definition_id": format!("sds-v1:{}", "b".repeat(64)), + })) + .unwrap(), + full_window_slide_ms: Some(slide), + materialization: asap_types::sds::StoredOutputId(7), + output_grouping: PhysicalGrouping::Reduce(vec![]), + item_labels: vec![], + window_ms: WINDOW, + pane_origin_ms: Some(0), + readout_lookback_ms: Some(WINDOW), + } + } + fn entry(nodes: Vec<(u64, QueryPlanNode)>, program: CompiledPhysicalDag) -> QueryPlanEntry { + QueryPlanEntry { + physical_dag: Some(serde_json::from_slice(&program.encode().unwrap()).unwrap()), + language: QueryLanguage::PromQl, + query_id: "mixed".into(), + canonical_query: "mixed".into(), + fixed_evaluation: None, + root: QueryNodeId(9), + nodes: nodes + .into_iter() + .map(|(id, node)| (QueryNodeId(id), node)) + .collect(), + instant: InstantExecution { + lookback_ms: WINDOW, + full_history: false, + cumulative_readout: false, + }, + fallback: FallbackPolicy::ExactBackend, + } + } + fn physical(inputs: Vec) -> QueryPlanNode { + QueryPlanNode::Physical { + source_nodes: inputs.clone(), + inputs: inputs.into_iter().map(QueryNodeId).collect(), + max_bytes: 1 << 20, + drop_metric_name: true, + } + } + fn mixed(slide: u64) -> QueryPlanEntry { + use planner_types::{post_asap::BinaryOperator, pre_asap::*}; + let window = Operator::series_window( + raw_schema(), + Some(AggIntent::Sum { col: None }), + WINDOW as i64, + 0, + None, + None, + ) + .unwrap(); + let add = Operator::series_binary( + raw_schema(), + stored_readout().schema(), + BinaryOperator { + checked_relative_division: false, + checked_finite_division: false, + kind: BinaryOpKind::Arithmetic(ArithmeticOpKind::Add), + vector_match: None, + }, + ) + .unwrap(); + let program = CompiledPhysicalDag::from_operators( + [ + (0, InputContract::bounded(raw_schema())), + (1, InputContract::bounded(stored_schema())), + ] + .into(), + [ + (2, (vec![0], window)), + (3, (vec![1], stored_readout())), + (4, (vec![2, 3], add)), + ] + .into(), + vec![4], + ) + .unwrap(); + let scan = QueryPlanNode::Logical { + operator: scan(), + inputs: vec![], + }; + let stored = QueryPlanNode::ReadMaterialization { + binding: binding(slide), + }; + entry( + vec![(0, scan), (1, stored), (9, physical(vec![0, 1]))], + program, + ) + } + /// Per-series stored Sum states for instances a and b. + fn stored(a: f64, b: f64) -> Batch { + let rows = [("a", a), ("b", b)] + .into_iter() + .map(|(instance, value)| { + let labels = BTreeMap::from([ + ("__name__".to_string(), "m".to_string()), + ("instance".to_string(), instance.to_string()), + ("job".to_string(), "api".to_string()), + ]); + let mut state = ExactAccumulator::new(sum_family(), false).unwrap(); + state.update(None, value, 0); + vec![ + Value::Utf8(encode_series_identity(&labels).unwrap().into()), + Value::Summary { + family: sum_family(), + state: Arc::new(state), + }, + ] + }) + .collect(); + Batch::try_new(stored_schema(), rows).unwrap() + } + type Store = BTreeMap NativeReadError>>; + /// Serve stored windows by end time, recording each requested window. + fn run( + entry: &QueryPlanEntry, + endpoint: &str, + max_lag: Option, + at: u64, + windows: &Store, + ) -> ( + Result< + ( + crate::query_engines::query_result::QueryResult, + super::super::super::logical_dag::ExecutionStats, + ), + EngineError, + >, + Vec, + ) { + let client = reqwest::Client::new(); + let mut requested = Vec::new(); + let result = super::super::super::logical_dag::native_values::execute_stored( + entry, + Some((&client, endpoint)), + max_lag, + at, + &mut |binding, window, schema, _| { + assert_eq!(binding.window_ms, WINDOW); + requested.push(window); + assert_eq!(window.end_ms - window.start_ms, WINDOW as i64); + match windows.get(&(window.end_ms as u64)) { + Some(Ok(batch)) => { + assert_eq!(batch.schema(), &schema); + Ok(batch.clone()) + } + Some(Err(error)) => Err(error()), + None => Err(NativeReadError::Unavailable("window incomplete".into())), + } + }, + ); + (result, requested) + } + fn values( + result: crate::query_engines::query_result::QueryResult, + ) -> BTreeMap { + let crate::query_engines::query_result::QueryResult::Vector(result) = result else { + panic!("expected an instant vector") + }; + assert_eq!(result.timestamp, AT_MS, "the result is at the query time"); + result + .values + .iter() + .map(|point| { + let keys = point.label_keys_override.clone().unwrap(); + let labels = keys + .into_iter() + .zip(point.labels.labels.iter().cloned()) + .collect::>(); + (labels["instance"].clone(), point.value) + }) + .collect() + } + async fn in_blocking(f: impl FnOnce() -> T + Send + 'static) -> T { + crate::query_engines::request::run(Default::default(), move |_| Ok(f())) + .await + .unwrap() + } + + // Within the default bound (one 60s slide), raw samples at t_q plus the + // newest complete stored window (lag 40s) equal the exact reference + // over that lagged stored state, and the lag is reported. + #[tokio::test] + async fn mixed_program_within_default_lag_matches_exact_reference() { + let (endpoint, seen, server) = + prometheus(200, matrix(), std::time::Duration::ZERO).await; + let (result, requested) = in_blocking(move || { + let windows = Store::from([(LATEST, Ok(stored(100., 200.)))]); + run(&mixed(SLIDE), &endpoint, None, AT_MS, &windows) + }) + .await; + let (result, stats) = result.unwrap(); + // Exact reference: sum_over_time of raw samples in (700s, 1000s] + // plus the stored window (660s, 960s] Sum. + assert_eq!( + values(result), + BTreeMap::from([ + ("a".into(), 5. + 1. + 3. + 100.), + ("b".into(), 10. + 50. + 40. + 20. + 30. + 200.), + ]) + ); + assert_eq!( + requested, + [HalfOpenTimeRange { + start_ms: (LATEST - WINDOW) as i64, + end_ms: LATEST as i64 + }] + ); + assert_eq!(stats.stored_input_lag_ms, Some(AT_MS - LATEST)); + assert_eq!( + seen.lock().unwrap()[0]["time"], + "1000.000", + "raw is read at t_q" + ); + server.abort(); + } + + // An incomplete newest window steps back one slide; that older window + // is admitted only when the configured bound covers its lag. + #[tokio::test] + async fn older_window_is_admitted_only_within_the_configured_bound() { + let (endpoint, _, server) = prometheus(200, matrix(), std::time::Duration::ZERO).await; + let older = LATEST - SLIDE; + let (admitted, beyond) = in_blocking(move || { + let windows = Store::from([(older, Ok(stored(1000., 2000.)))]); + ( + run( + &mixed(SLIDE), + &endpoint, + Some(AT_MS - older), + AT_MS, + &windows, + ), + run( + &mixed(SLIDE), + &endpoint, + Some(AT_MS - older - 1), + AT_MS, + &windows, + ), + ) + }) + .await; + let (result, stats) = admitted.0.unwrap(); + assert_eq!(values(result)["a"], 5. + 1. + 3. + 1000.); + assert_eq!(stats.stored_input_lag_ms, Some(AT_MS - older)); + assert_eq!(admitted.1.len(), 2, "newest window, then one slide older"); + assert!( + matches!(&beyond.0, Err(EngineError::CapabilityMiss { detail, .. }) + if detail.contains("lag 100000ms exceeds bound 99999ms")), + "{:?}", + beyond.0.as_ref().err() + ); + server.abort(); + } + + // Beyond the default bound the query is a capability miss, which takes + // the exact fallback, before any raw read; a stored resource failure + // stays terminal instead of stepping to an older window. + #[tokio::test] + async fn stale_stored_input_falls_back_before_reading_raw_series() { + let (endpoint, seen, server) = + prometheus(200, matrix(), std::time::Duration::ZERO).await; + let (stale, failed) = in_blocking(move || { + let stale = Store::from([(LATEST - SLIDE, Ok(stored(1., 2.)))]); + let failed: Store = BTreeMap::from([( + LATEST, + Err( + (|| NativeReadError::Physical(asap_physical_operators::Error::MemoryLimit)) + as fn() -> NativeReadError, + ), + )]); + ( + run(&mixed(SLIDE), &endpoint, None, AT_MS, &stale), + run(&mixed(SLIDE), &endpoint, None, AT_MS, &failed), + ) + }) + .await; + assert!( + matches!(&stale.0, Err(EngineError::CapabilityMiss { detail, .. }) + if detail.contains("exceeds bound 60000ms")), + "{:?}", + stale.0.as_ref().err() + ); + assert!( + seen.lock().unwrap().is_empty(), + "no raw read after a stale stored input" + ); + assert!(matches!( + failed.0, + Err(EngineError::Physical( + asap_physical_operators::Error::MemoryLimit + )) + )); + assert_eq!(failed.1.len(), 1); + server.abort(); + } + + // Raw inputs mix only with stored native batches, whose lag is checked. + #[test] + fn raw_inputs_do_not_mix_with_other_stored_readouts() { + let mut entry = mixed(SLIDE); + let stored = entry.nodes.insert( + QueryNodeId(1), + QueryPlanNode::ExactReadout { + input: QueryNodeId(2), + readout: asap_types::query_plan::ExactReadout::Rate, + }, + ); + entry.nodes.insert(QueryNodeId(2), stored.unwrap()); + let Err(error) = entry.recover_vector_physical_dag() else { + panic!("raw input mixed with an exact readout") + }; + assert!( + error + .to_string() + .contains("mix only with stored native batches"), + "{error}" + ); + assert!(mixed(SLIDE).recover_vector_physical_dag().is_ok()); + } + + // A stored-only program still reads exactly the window ending at the + // evaluation time and reports no lag. + #[tokio::test] + async fn stored_only_program_reads_the_window_ending_at_query_time() { + let program = CompiledPhysicalDag::from_operators( + [(1, InputContract::bounded(stored_schema()))].into(), + [(3, (vec![1], stored_readout()))].into(), + vec![3], + ) + .unwrap(); + let stored_only = entry( + vec![ + ( + 1, + QueryPlanNode::ReadMaterialization { + binding: binding(SLIDE), + }, + ), + (9, physical(vec![1])), + ], + program, + ); + let (result, requested) = in_blocking(move || { + let windows = + Store::from([(AT_MS, Ok(stored(7., 8.))), (LATEST, Ok(stored(1., 2.)))]); + run(&stored_only, "http://unused", None, AT_MS, &windows) + }) + .await; + let (result, stats) = result.unwrap(); + assert_eq!(values(result)["a"], 7.); + assert_eq!( + requested, + [HalfOpenTimeRange { + start_ms: AT - WINDOW as i64, + end_ms: AT + }] + ); + assert_eq!(stats.stored_input_lag_ms, None); + } + } } From 0adc35a3f827133875f33b25bc094e7d1d3e481c Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 12:36:13 +0000 Subject: [PATCH 3/3] feat(http): report stored-input lag as a response header Co-Authored-By: Claude Opus 5.5 --- data_plane/src/drivers/query/servers/http.rs | 39 ++++++++++++++++++++ 1 file changed, 39 insertions(+) diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index 7a0c94dc1..610f57a5b 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -1488,6 +1488,23 @@ async fn process_via_router( } } +/// Remove the engine's stored-input lag note from the client warnings. +fn extract_stored_input_lag(value: &mut serde_json::Value) -> Option { + let warnings = value.get_mut("warnings")?.as_array_mut()?; + let mut lag = None; + warnings.retain(|warning| { + let Some(text) = warning + .as_str() + .and_then(|text| text.strip_prefix("asap_stored_input_lag_ms:")) + else { + return true; + }; + lag = text.parse().ok(); + false + }); + lag +} + fn extract_logical_provenance( value: &mut serde_json::Value, ) -> Option> { @@ -1622,6 +1639,12 @@ async fn annotate_data_source(response: Response, data_source_id: &'static str) return Response::from_parts(parts, axum::body::Body::from(bytes)); } if data_source_id == "asap_query" { + if let Some(lag) = extract_stored_input_lag(&mut value) { + parts.headers.insert( + "x-asap-stored-input-lag-ms", + axum::http::HeaderValue::from(lag), + ); + } if let Some(provenance) = extract_logical_provenance(&mut value) { let (route, detail) = match provenance { Ok((raw, summary, memo, remote, _legacy_indexes, rpcs, branches)) => { @@ -6477,6 +6500,22 @@ mod logical_provenance_tests { assert_eq!(value["warnings"], serde_json::json!(["partial data"])); } + // The stored-input lag of a mixed program becomes response metadata, not a client warning. + #[tokio::test] + async fn stored_input_lag_is_a_response_header() { + let response = Json(serde_json::json!({"status":"success", "warnings":[ + "asap_stored_input_lag_ms:40000", "asap_logical_stats:raw=0,summary=2,memo_hits=0,remote=0,remote_rpcs=0,remote_branches=0" + ], "data":{"resultType":"vector", "result":[]}})).into_response(); + let response = annotate_data_source(response, "asap_query").await; + assert_eq!(response.headers()["x-asap-stored-input-lag-ms"], "40000"); + assert_eq!(response.headers()["x-asap-execution"], "warm"); + let bytes = axum::body::to_bytes(response.into_body(), 4096) + .await + .unwrap(); + let value: serde_json::Value = serde_json::from_slice(&bytes).unwrap(); + assert!(value["warnings"].as_array().unwrap().is_empty()); + } + #[test] fn contradictory_provenance_is_not_warm() { // Any observed local raw branch invalidates a deployed plan.