Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
56 changes: 33 additions & 23 deletions data_plane/src/drivers/ingest/prometheus_remote_write.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1736,7 +1736,7 @@ impl PrometheusRemoteWriteReceiver {
now_ms: u64,
) -> Result<u64, crate::precompute_engine::revisions::RevisionError> {
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
Expand Down Expand Up @@ -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::<Vec<_>>();
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::<asap_physical_operators::Error>() {
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(())
Expand All @@ -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<dyn asap_summary_state::summary_kernels::factory::AccumulatorUpdater>> = 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<str>,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<str> = 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::<usize>();
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.
Expand All @@ -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::<usize>() > 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<dyn crate::storage_engines::types::AggregateCore>> = windows.into_iter().map(|(w,updater)|(w,Arc::from(updater.into_accumulator()))).collect();
let mut states: BTreeMap<_, Arc<dyn crate::storage_engines::types::AggregateCore>> = 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())? });
}
Expand Down
144 changes: 89 additions & 55 deletions data_plane/src/precompute_engine/raw_dag.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<dyn std::error::Error + Send + Sync>;

/// A validated executable projection; semantics come from the installed node.
/// The retained node ID makes failures attributable to the selected DAG.
#[derive(Debug, Clone)]
Expand All @@ -27,7 +31,6 @@ pub struct RawDagProgram {
pub input: SummaryUpdate,
pub grouping: GroupingStrategy,
pub reduction: planner_types::pre_asap::Reduction,
projected_column: Option<String>,
/// 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]>,
Expand Down Expand Up @@ -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
Expand All @@ -222,29 +223,83 @@ 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>(
&self,
samples: impl IntoIterator<Item = (&'a str, i64, f64)>,
pane: (i64, i64),
max_bytes: usize,
) -> Result<Option<Box<dyn AggregateCore>>, String> {
use futures::StreamExt;
) -> Result<Option<Box<dyn AggregateCore>>, 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<Item = (&'a str, i64, f64)>,
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<Box<dyn AggregateCore>, String> {
Ok(self.heap_updater()?.take_accumulator())
}

fn execute<'a>(
&self,
samples: impl IntoIterator<Item = (&'a str, i64, f64)>,
pane: (i64, i64),
max_bytes: usize,
) -> Result<Vec<Vec<Value>>, BuildError> {
use futures::StreamExt;
let schema = precompute::raw_sample_schema();
let rows = samples
.into_iter()
Expand Down Expand Up @@ -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()),
)?;
Expand All @@ -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<Box<dyn AccumulatorUpdater>, 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 {
Expand All @@ -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<f64, String> {
fn heap_updater(&self) -> Result<Box<dyn AccumulatorUpdater>, String> {
create_planner_accumulator(&self.family, &self.input, &self.grouping)
}

fn validate_sample(&self, updater: &dyn AccumulatorUpdater, value: f64) -> Result<f64, String> {
let weight = match &self.input.weight {
SummaryInputExpr::Constant(c) => *c,
// The worker retains one previous value per series across pane rotation.
Expand All @@ -347,7 +381,7 @@ impl RawDagProgram {
Ok(weight)
}

pub fn apply(
fn apply(
&self,
updater: &mut dyn AccumulatorUpdater,
series: &str,
Expand Down
11 changes: 6 additions & 5 deletions data_plane/src/precompute_engine/revisions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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));
}
Expand Down
Loading
Loading