diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index 6acc03e5c..ff117d7b0 100644 --- a/data_plane/src/drivers/ingest/prometheus_remote_write.rs +++ b/data_plane/src/drivers/ingest/prometheus_remote_write.rs @@ -1736,7 +1736,7 @@ impl PrometheusRemoteWriteReceiver { now_ms: u64, ) -> Result { use crate::precompute_engine::revisions::{encode_state, InputSample, RevisionRecord}; - use crate::precompute_engine::{raw_dag::RawDagProgram, window_manager::WindowManager}; + use crate::precompute_engine::window_manager::WindowManager; use asap_types::sds::StoredOutputId; let (plan, store) = runtime.installed()?; let generation = plan @@ -1825,15 +1825,20 @@ impl PrometheusRemoteWriteReceiver { } for config in plan.precompute_plan.materializations.iter().filter(|c| c.derived_input.is_none()) { let filter = compile_spatial_filter(&config.spatial_filter_normalized)?; - let program = RawDagProgram::from_plan(&plan.precompute_plan, config)?; - let updater = program.updater()?; - for sample in samples { + let program = plan.installed_precompute_plan.raw_programs.get(&config.policy_fp_u64()).ok_or("revision raw program missing")?; + let mut matching = samples.iter().filter_map(|sample| { let labels = sample.labels.iter().map(|(k,v)|(k.clone(),v.clone())).collect(); - if sample.metric == config.metric && filter.matches(&labels) { - if let Some(value) = sample.value { - program.validate_sample(updater.as_ref(), value).map_err(|_| crate::precompute_engine::revisions::AdmissionRejected("sample is outside the bound summary kernel's input domain"))?; - } - } + (sample.metric == config.metric && filter.matches(&labels)).then_some(())?; + Some((sample.series.as_str(), sample.timestamp_ms, sample.value?)) + }).collect::>(); + matching.sort_by(|a, b| (a.0, a.1).cmp(&(b.0, b.1))); + // Reject the whole request when the Planner graph cannot admit its + // values; resource and setup failures are not input errors. + if let Err(error) = program.validate(matching, runtime.policy.max_checkpoint_bytes) { + return Err(match error.downcast_ref::() { + Some(asap_physical_operators::Error::MemoryLimit) | None => error, + Some(_) => crate::precompute_engine::revisions::AdmissionRejected("sample is outside the bound summary kernel's input domain").into(), + }); } } Ok(()) @@ -1856,19 +1861,18 @@ impl PrometheusRemoteWriteReceiver { let WorkerMessage::GroupSamples { sid, policy_fp, group_key, mut samples, .. } = message else { return Err("revision routing returned non-sample input".into()); }; let config = plan.precompute_plan.materializations.iter().find(|c| c.policy_fingerprint() == policy_fp).ok_or("revision raw configuration missing")?; if committed.contains_key(&policy_fp.0) { continue; } - let program = RawDagProgram::from_plan(&plan.precompute_plan, config)?; + let program = plan.installed_precompute_plan.raw_programs.get(&config.policy_fp_u64()).ok_or("revision raw program missing")?; let manager = WindowManager::with_layout(config.window_size, config.slide_interval, config.pane_origin_ms, &config.window_layout); - let mut windows: BTreeMap<(u64,u64), Box> = BTreeMap::new(); + // Each stored window's admitted samples, built by the Planner graph below. + // Overlapping windows share each sample's series key. + let mut windows: BTreeMap<(u64,u64), Vec<(Arc,i64,f64)>> = BTreeMap::new(); samples.sort_by(|a,b| (&a.0,a.1).cmp(&(&b.0,b.1))); for (series,time,value) in samples { + let series: Arc = Arc::from(series); for start in manager.stored_bucket_starts(time) { let (start,end) = manager.stored_bucket_bounds(start); if start < 0 || end < 0 || end as u64 > captured_at_ms || (start as u64) < captured_at_ms.saturating_sub(retain.saturating_sub(max_window)) { continue; } - let key = (start as u64,end as u64); - if let std::collections::btree_map::Entry::Vacant(entry) = windows.entry(key) { entry.insert(program.updater()?); } - program.apply(windows.get_mut(&key).unwrap().as_mut(), &series, value, time)?; - let bytes = windows.values().map(|u|u.memory_usage_bytes()).sum::(); - if bytes > runtime.policy.max_checkpoint_bytes { return Err(Box::new(asap_physical_operators::Error::MemoryLimit)); } + windows.entry((start as u64,end as u64)).or_default().push((Arc::clone(&series), time, value)); } } // This captured revision contains every accepted local input. @@ -1886,19 +1890,25 @@ impl PrometheusRemoteWriteReceiver { if width == 0 { return Err("counter pane width is zero".into()); } while start < last { let end = start.checked_add(width).ok_or("counter pane timestamp overflow")?; - if let std::collections::btree_map::Entry::Vacant(entry) = windows.entry((start, end)) { - entry.insert(program.updater()?); - if windows.values().map(|u| u.memory_usage_bytes()).sum::() > runtime.policy.max_checkpoint_bytes { - return Err(Box::new(asap_physical_operators::Error::MemoryLimit)); - } - } + windows.entry((start, end)).or_default(); start = end; } } } let reference = plan.installed_precompute_plan.stored_output_reference(policy_fp.into()).ok_or("revision raw binding missing")?; let group = group_key.as_population_labels(); - let states: BTreeMap<_, Arc> = windows.into_iter().map(|(w,updater)|(w,Arc::from(updater.into_accumulator()))).collect(); + let mut states: BTreeMap<_, Arc> = BTreeMap::new(); + for (window, samples) in windows { + let state = if samples.is_empty() { + // A known-empty counter pane of this captured revision. + program.empty_state()? + } else { + let bounds = (i64::try_from(window.0)?, i64::try_from(window.1)?); + program.build(samples.iter().map(|(s,t,v)|(s.as_ref(),*t,*v)), bounds, runtime.policy.max_checkpoint_bytes)? + .ok_or("revision window admitted no population")? + }; + states.insert(window, Arc::from(state)); + } for ((start,end),state) in &states { records.get_mut(&policy_fp.into()).unwrap().push(RevisionRecord { reference: reference.clone(), group: group.clone(), start_ms:*start, end_ms:*end, payload:encode_state(Arc::clone(state),program.family.clone())? }); } diff --git a/data_plane/src/precompute_engine/raw_dag.rs b/data_plane/src/precompute_engine/raw_dag.rs index 2180b669c..1f10204cc 100644 --- a/data_plane/src/precompute_engine/raw_dag.rs +++ b/data_plane/src/precompute_engine/raw_dag.rs @@ -18,6 +18,10 @@ use planner_types::post_asap::{ use planner_types::pre_asap::{ColumnRef, QueryExpr, Source}; use std::collections::HashMap; +/// Planner execution errors keep their type (e.g. `MemoryLimit`); setup and +/// binding failures are messages. +pub type BuildError = Box; + /// A validated executable projection; semantics come from the installed node. /// The retained node ID makes failures attributable to the selected DAG. #[derive(Debug, Clone)] @@ -27,7 +31,6 @@ pub struct RawDagProgram { pub input: SummaryUpdate, pub grouping: GroupingStrategy, pub reduction: planner_types::pre_asap::Reduction, - projected_column: Option, /// Planner's encoded precompute graph from the raw sample boundary to this /// output. Decoded graphs are not `Send`, so each execution decodes it. program: std::sync::Arc<[u8]>, @@ -197,12 +200,10 @@ impl RawDagProgram { input: input.clone(), grouping: grouping.clone(), reduction: reduction.clone(), - projected_column: config - .effective_value_projection() - .column() - .map(str::to_owned), }; - program.validate()?; + if program.heap() { + program.validate_heap_update()?; + } if let Some(old) = &selected { if old.family != program.family || old.input != program.input @@ -222,6 +223,16 @@ impl RawDagProgram { selected.ok_or_else(|| "raw materialization has no selected post-ASAP DAG producer".into()) } + /// Heap readout decodes only the backend heap kernel, not Planner's weighted + /// frequency state, so heaps stay on that kernel until it does. + fn heap(&self) -> bool { + matches!(&self.family, SummaryFamilyType::Sketch(kind, _) if matches!( + kind.algorithm(), + planner_types::post_asap::SketchAlgorithm::CmsWithHeap + | planner_types::post_asap::SketchAlgorithm::CountSketchWithHeap + )) + } + /// Execute the Planner graph over one pane's samples as one typed batch. /// Returns `None` when the graph admits no population from them. pub fn build<'a>( @@ -229,22 +240,66 @@ impl RawDagProgram { samples: impl IntoIterator, pane: (i64, i64), max_bytes: usize, - ) -> Result>, String> { - use futures::StreamExt; + ) -> Result>, BuildError> { let samples = pane_batch(samples); - // Stored heap readout decodes only the backend heap kernel, not - // Planner's weighted frequency state, so heaps stay on that kernel. - if matches!(&self.family, SummaryFamilyType::Sketch(kind, _) if matches!( - kind.algorithm(), - planner_types::post_asap::SketchAlgorithm::CmsWithHeap - | planner_types::post_asap::SketchAlgorithm::CountSketchWithHeap - )) { - let mut updater = self.updater()?; + if self.heap() { + let mut updater = self.heap_updater()?; for (series, time, value) in samples { self.apply(&mut *updater, series, value, time)?; } return Ok(Some(updater.take_accumulator())); } + // The router assigns one population per group, so one pane yields + // at most one state. + match self.execute(samples, pane, max_bytes)?.as_slice() { + [] => Ok(None), + [row] => match row.as_slice() { + [_, _, Value::Summary { state, .. }] => Ok(Some( + asap_summary_state::physical::from_physical(state.as_ref())?, + )), + _ => Err("raw precompute output is not a population state".into()), + }, + _ => Err("one routed group produced several populations".into()), + } + } + + /// Check that the Planner graph admits every sample, without storing a result. + pub fn validate<'a>( + &self, + samples: impl IntoIterator, + max_bytes: usize, + ) -> Result<(), BuildError> { + let samples = pane_batch(samples); + if self.heap() { + let updater = self.heap_updater()?; + for (_, _, value) in samples { + self.validate_sample(updater.as_ref(), value)?; + } + return Ok(()); + } + if samples.is_empty() { + return Ok(()); + } + let bounds = samples + .iter() + .fold((i64::MAX, i64::MIN), |(lo, hi), (_, time, _)| { + (lo.min(*time), hi.max(time.saturating_add(1))) + }); + self.execute(samples, bounds, max_bytes).map(|_| ()) + } + + /// The family's empty state, for a pane known to have no samples. + pub fn empty_state(&self) -> Result, String> { + Ok(self.heap_updater()?.take_accumulator()) + } + + fn execute<'a>( + &self, + samples: impl IntoIterator, + pane: (i64, i64), + max_bytes: usize, + ) -> Result>, BuildError> { + use futures::StreamExt; let schema = precompute::raw_sample_schema(); let rows = samples .into_iter() @@ -272,7 +327,7 @@ impl RawDagProgram { }, ) .map_err(|e| e.to_string())?; - let rows = futures::executor::block_on(async { + Ok(futures::executor::block_on(async { let mut stream = graph.execute(program.roots(), context)?.pop().ok_or( asap_physical_operators::Error::Invalid("missing output".into()), )?; @@ -281,35 +336,16 @@ impl RawDagProgram { rows.extend(batch?.rows().iter().cloned()); } Ok::<_, asap_physical_operators::Error>(rows) - }) - .map_err(|e| e.to_string())?; - // The router assigns one population per group, so one pane yields - // at most one state. - match rows.as_slice() { - [] => Ok(None), - [row] => match row.as_slice() { - [_, _, Value::Summary { state, .. }] => { - asap_summary_state::physical::from_physical(state.as_ref()) - .map(Some) - .map_err(|e| e.to_string()) - } - _ => Err("raw precompute output is not a population state".into()), - }, - _ => Err("one routed group produced several populations".into()), - } + })?) } - pub fn updater(&self) -> Result, String> { - create_planner_accumulator(&self.family, &self.input, &self.grouping) - } - - fn validate(&self) -> Result<(), String> { - match &self.input.weight { - SummaryInputExpr::Column(ColumnRef::SampleValue) | SummaryInputExpr::Constant(_) => {} - SummaryInputExpr::Column( - ColumnRef::Named(name) | ColumnRef::Qualified { name, .. }, - ) if self.projected_column.as_ref() == Some(name) => {} - _ => return Err("raw DAG weight expression is unsupported".into()), + /// Heaps still use the kernel interpreter, so install only updates it evaluates. + fn validate_heap_update(&self) -> Result<(), String> { + if !matches!( + &self.input.weight, + SummaryInputExpr::Column(ColumnRef::SampleValue) | SummaryInputExpr::Constant(_) + ) { + return Err("raw heap weight expression is unsupported".into()); } fn item(expr: &SummaryInputExpr) -> bool { match expr { @@ -321,19 +357,17 @@ impl RawDagProgram { _ => false, } } - if self.input.item.as_ref().is_some_and(|e| !item(e)) { - return Err("raw DAG item expression is unsupported".into()); + if !self.input.item.as_ref().is_some_and(item) { + return Err("raw heap item expression is unsupported".into()); } - self.updater().map(|_| ()) + self.heap_updater().map(|_| ()) } - /// Validate admission with the same weight semantics used during execution, - /// without mutating an accumulator or accepting part of a request. - pub fn validate_sample( - &self, - updater: &dyn AccumulatorUpdater, - value: f64, - ) -> Result { + fn heap_updater(&self) -> Result, String> { + create_planner_accumulator(&self.family, &self.input, &self.grouping) + } + + fn validate_sample(&self, updater: &dyn AccumulatorUpdater, value: f64) -> Result { let weight = match &self.input.weight { SummaryInputExpr::Constant(c) => *c, // The worker retains one previous value per series across pane rotation. @@ -347,7 +381,7 @@ impl RawDagProgram { Ok(weight) } - pub fn apply( + fn apply( &self, updater: &mut dyn AccumulatorUpdater, series: &str, diff --git a/data_plane/src/precompute_engine/revisions.rs b/data_plane/src/precompute_engine/revisions.rs index 21c9c55ae..3b2521329 100644 --- a/data_plane/src/precompute_engine/revisions.rs +++ b/data_plane/src/precompute_engine/revisions.rs @@ -592,11 +592,12 @@ impl RevisionRuntime { .iter() .filter(|c| c.derived_input.is_none()) { - let program = super::raw_dag::RawDagProgram::from_plan(&plan.precompute_plan, config)?; - let bytes = encode_state( - Arc::from(program.updater()?.into_accumulator()), - program.family, - )?; + let program = plan + .installed_precompute_plan + .raw_programs + .get(&config.policy_fp_u64()) + .ok_or("revision raw program missing")?; + let bytes = encode_state(Arc::from(program.empty_state()?), program.family.clone())?; if bytes.len() > self.policy.max_checkpoint_bytes { return Err(Box::new(asap_physical_operators::Error::MemoryLimit)); } diff --git a/data_plane/src/precompute_engine/worker.rs b/data_plane/src/precompute_engine/worker.rs index 814ed12d0..29a3b8bb5 100644 --- a/data_plane/src/precompute_engine/worker.rs +++ b/data_plane/src/precompute_engine/worker.rs @@ -1596,13 +1596,15 @@ fn build_pane( pane: (i64, i64), ) -> Result>, String> { if let Some(program) = program { - return program.build( - samples - .iter() - .map(|(series, time, value)| (series.as_ref(), *time, *value)), - pane, - asap_physical_operators::runtime::Limits::default().max_bytes, - ); + return program + .build( + samples + .iter() + .map(|(series, time, value)| (series.as_ref(), *time, *value)), + pane, + asap_physical_operators::runtime::Limits::default().max_bytes, + ) + .map_err(|e| e.to_string()); } #[cfg(test)] { @@ -4658,8 +4660,17 @@ mod dag_execution_tests { plans } - // Live raw ingest builds each pane with the installed Planner graph, and - // every stored state equals feeding the selected kernel sample by sample. + /// The pre-Planner interpreter's unkeyed update: a constant weight or the + /// sample value; unit-frequency summaries observe the sample value. + fn reference_weight(input: &planner_types::post_asap::SummaryUpdate, value: f64) -> f64 { + match (&input.item, &input.weight) { + (None, planner_types::post_asap::SummaryInputExpr::Constant(weight)) => *weight, + _ => value, + } + } + + // Live raw ingest and backfill build each pane with the installed Planner + // graph, and every stored state equals feeding the kernel sample by sample. #[test] fn live_panes_execute_planner_dag_with_identical_states() { let mut outputs = 0; @@ -4747,6 +4758,7 @@ mod dag_execution_tests { .and_then(serde_json::Value::as_bool) .unwrap_or(false); let mut expected = BTreeMap::new(); + let mut backfill = BTreeMap::new(); for (key, samples) in &groups { let group_key = Arc::new(GroupKey::new( key.iter().map(|(k, v)| (k.as_str(), v.as_str())), @@ -4761,10 +4773,27 @@ mod dag_execution_tests { start as u64, end as u64, )) - .or_insert_with(|| program.updater().unwrap()); - program - .apply(&mut **updater, series, *value, *time) - .unwrap(); + .or_insert_with(|| { + asap_summary_state::factory::create_planner_accumulator( + &program.family, + &program.input, + &program.grouping, + ) + .unwrap() + }); + updater.update_single(reference_weight(&program.input, *value), *time); + backfill + .entry(( + group_key.values().labels.join(";"), + start as u64, + end as u64, + )) + .or_insert_with(Vec::new) + .push(crate::storage_engines::sketch_db::backfill::RawSample { + labels: series.clone(), + timestamp_ms: *time, + value: *value, + }); } } } @@ -4810,6 +4839,25 @@ mod dag_execution_tests { config.metric ); assert_eq!(actual, expected, "{:?}", config.aggregation_type); + // Backfill builds each window with the same installed graph. + let backfilled = backfill + .into_iter() + .map(|((key, start, end), samples)| { + let state = crate::storage_engines::sketch_db::build_dag_accumulator( + &program, + &samples, + (start, end), + ) + .unwrap() + .unwrap(); + ((key, start, end), state.serialize_to_bytes()) + }) + .collect::>(); + assert_eq!( + backfilled, expected, + "backfill {:?}", + config.aggregation_type + ); outputs += 1; families.insert( format!("{:?}", config.accumulator_spec().unwrap().family) @@ -4986,6 +5034,37 @@ mod dag_execution_tests { } } + // Revision admission checks samples with the installed Planner graph, and + // keeps a Planner memory limit distinguishable from invalid input. + #[test] + fn admission_validates_with_the_planner_graph() { + let plan = plan("sum_over_time(asap_demo_gauge[5s])"); + let config = plan.precompute_plan.materializations[0].clone(); + let program = InstalledPrecomputePlan::from_precompute_plan(plan.precompute_plan) + .unwrap() + .raw_programs[&config.policy_fp_u64()] + .clone(); + let series = config.metric.as_str(); + assert!(program + .validate([(series, 1000, 1.0), (series, 2000, 2.0)], 1 << 20) + .is_ok()); + let rejected = program + .validate([(series, 1000, f64::INFINITY)], 1 << 20) + .unwrap_err(); + assert!(!matches!( + rejected.downcast_ref::(), + Some(asap_physical_operators::Error::MemoryLimit) + )); + // A resource limit stays typed, so admission does not report it as bad input. + let limited = program + .validate((0..64).map(|t| (series, t * 10, 1.0)), 1) + .unwrap_err(); + assert!(matches!( + limited.downcast_ref::(), + Some(asap_physical_operators::Error::MemoryLimit) + )); + } + // A raw output installs only with its Planner-compiled precompute graph. #[test] fn raw_output_requires_its_planner_precompute_graph() { diff --git a/data_plane/src/storage_engines/sketch_db/backfill/processor.rs b/data_plane/src/storage_engines/sketch_db/backfill/processor.rs index 12f391aaa..d408a54db 100644 --- a/data_plane/src/storage_engines/sketch_db/backfill/processor.rs +++ b/data_plane/src/storage_engines/sketch_db/backfill/processor.rs @@ -290,7 +290,12 @@ impl WindowProcessor for BackfillWindowProcessor { for (sid, bucket) in by_bucket { let SidBucket { group_key, samples } = bucket; let accumulator = if let Some(program) = &program { - super::window_builder::build_dag_accumulator(program, &samples)? + let Some(accumulator) = + super::window_builder::build_dag_accumulator(program, &samples, window_range)? + else { + continue; + }; + accumulator } else { #[cfg(test)] { diff --git a/data_plane/src/storage_engines/sketch_db/backfill/window_builder.rs b/data_plane/src/storage_engines/sketch_db/backfill/window_builder.rs index 4944be0ef..54db73627 100644 --- a/data_plane/src/storage_engines/sketch_db/backfill/window_builder.rs +++ b/data_plane/src/storage_engines/sketch_db/backfill/window_builder.rs @@ -161,21 +161,25 @@ mod tests { } } -/// Backfill uses the same selected DAG producer and update expressions as live input. +/// Backfill builds a window with the same installed Planner graph as live +/// input, over the window's samples in ingestion order. `None` means the graph +/// admitted no population from them. pub fn build_dag_accumulator( program: &crate::precompute_engine::raw_dag::RawDagProgram, samples: &[RawSample], -) -> Result, String> { - let mut updater = program.updater()?; - // A selected program never converts samples to counter deltas; rate is an - // explicit upstream operator, so each sample is applied as it arrives. - for sample in samples { - program.apply( - &mut *updater, - &sample.labels, - sample.value, - sample.timestamp_ms, - )?; - } - Ok(updater.take_accumulator()) + window: (u64, u64), +) -> Result>, String> { + let window = ( + i64::try_from(window.0).map_err(|_| "backfill window start overflow")?, + i64::try_from(window.1).map_err(|_| "backfill window end overflow")?, + ); + program + .build( + samples + .iter() + .map(|sample| (sample.labels.as_str(), sample.timestamp_ms, sample.value)), + window, + asap_physical_operators::runtime::Limits::default().max_bytes, + ) + .map_err(|e| e.to_string()) }