diff --git a/Cargo.lock b/Cargo.lock index f75e5382b..ea23af070 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -364,7 +364,7 @@ dependencies = [ [[package]] name = "asap-aware-mapping" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=9a13ae43eecc8487fe6a25a98f6ec09399d2073e#9a13ae43eecc8487fe6a25a98f6ec09399d2073e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=ca422df610009bdcc4f05b9becf4b2f5ee0c56b2#ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" dependencies = [ "asap-types", "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib)", @@ -376,7 +376,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=9a13ae43eecc8487fe6a25a98f6ec09399d2073e#9a13ae43eecc8487fe6a25a98f6ec09399d2073e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=ca422df610009bdcc4f05b9becf4b2f5ee0c56b2#ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" dependencies = [ "asap-types", "promql-parser 0.10.0 (git+https://github.com/ProjectASAP/promql-parser?rev=9fede7eecca923c9882fe256484d00d37f8706cb)", @@ -385,7 +385,7 @@ dependencies = [ [[package]] name = "asap-frontend-sql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=9a13ae43eecc8487fe6a25a98f6ec09399d2073e#9a13ae43eecc8487fe6a25a98f6ec09399d2073e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=ca422df610009bdcc4f05b9becf4b2f5ee0c56b2#ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -396,7 +396,7 @@ dependencies = [ [[package]] name = "asap-physical-operators" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=9a13ae43eecc8487fe6a25a98f6ec09399d2073e#9a13ae43eecc8487fe6a25a98f6ec09399d2073e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=ca422df610009bdcc4f05b9becf4b2f5ee0c56b2#ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" dependencies = [ "asap-types", "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib?rev=5f03ccbd798ed5fec62bdd839bcb331123cab369)", @@ -410,12 +410,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=9a13ae43eecc8487fe6a25a98f6ec09399d2073e#9a13ae43eecc8487fe6a25a98f6ec09399d2073e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=ca422df610009bdcc4f05b9becf4b2f5ee0c56b2#ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=9a13ae43eecc8487fe6a25a98f6ec09399d2073e#9a13ae43eecc8487fe6a25a98f6ec09399d2073e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=ca422df610009bdcc4f05b9becf4b2f5ee0c56b2#ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" dependencies = [ "serde", "serde_json", diff --git a/Cargo.toml b/Cargo.toml index a416f3a41..63cf197d0 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -16,10 +16,10 @@ version = "0.1.0" [workspace.dependencies] # Keep Planner frontends, selection, and IR on the same immutable revision. # Alias upstream asap-types because this workspace also defines asap_types. -planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" } -asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" } -asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" } -asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" } +planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" } +asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" } +asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" } +asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" } # Shared external deps (used by 2+ crates) serde = { version = "1.0", features = ["derive"] } @@ -39,9 +39,10 @@ arc-swap = "1.7" reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] } # Internal crates -asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" } +asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" } asap_sketch_codec = { path = "crates/asap_sketch_codec" } asap_summary_state = { path = "crates/asap_summary_state" } asap_types = { path = "crates/asap_types" } asap_otel_proto = { path = "crates/asap_otel_proto" } indexmap = { version = "2.0", features = ["serde"] } + diff --git a/control_plane/src/physical/executable_binding.rs b/control_plane/src/physical/executable_binding.rs index 7639251ea..cd046258b 100644 --- a/control_plane/src/physical/executable_binding.rs +++ b/control_plane/src/physical/executable_binding.rs @@ -100,6 +100,7 @@ pub(super) fn compile_precompute_programs( use asap_types::executable_plan::BackendNodeBinding; let dag = installed.document.decode()?; let mut groups = std::collections::BTreeMap::<_, Vec>::new(); + let mut raw = Vec::new(); for sink in &installed.binding.precompute_sinks { if installed.native_programs.contains_key(sink) { continue; @@ -114,6 +115,24 @@ pub(super) fn compile_precompute_programs( .find(|c| c.policy_fingerprint() == stored_output.fingerprint()) .ok_or("precompute output configuration missing")?; let Some(derived) = &config.derived_input else { + // Raw ingest: Planner compiles the summary over its raw sample boundary. + let raw_inputs = dag + .edges + .iter() + .filter(|edge| edge.consumer == *sink) + .map(|edge| u64::from(edge.producer.0)) + .collect::>(); + let program = asap_physical_operators::physical_planner::precompute::compile( + &dag, + &raw_inputs, + &[u64::from(sink.0)], + ) + .map_err(|e| format!("raw precompute output {}: {e}", sink.0))?; + raw.push(( + *sink, + serde_json::from_slice(&program.encode().map_err(|e| e.to_string())?) + .map_err(|e| e.to_string())?, + )); continue; }; let frontiers = installed.binding.nodes.iter().filter_map(|(id, binding)| { @@ -133,6 +152,7 @@ pub(super) fn compile_precompute_programs( .or_default() .push(u64::from(sink.0)); } + installed.native_programs.extend(raw); for ((frontiers, _, _, _, _), roots) in groups { let program = asap_physical_operators::physical_planner::precompute::compile( &dag, &frontiers, &roots, diff --git a/control_plane/src/physical/plan_dot.rs b/control_plane/src/physical/plan_dot.rs index 2a7c5e980..ba2033fba 100644 --- a/control_plane/src/physical/plan_dot.rs +++ b/control_plane/src/physical/plan_dot.rs @@ -71,7 +71,11 @@ pub fn render(plan: &CompiledPhysicalPlan) -> String { render_native( &mut dot, &format!("maintenance_{dag_index}_{}", sink.0), - "Planner maintenance operators", + if installed.reads_raw_samples(*sink) { + "Planner raw precompute operators" + } else { + "Planner maintenance operators" + }, program, ); } diff --git a/control_plane/src/physical/workload_cost.rs b/control_plane/src/physical/workload_cost.rs index 32968be90..e5740894d 100644 --- a/control_plane/src/physical/workload_cost.rs +++ b/control_plane/src/physical/workload_cost.rs @@ -215,6 +215,10 @@ pub fn manifest( .iter() .find(|m| m.policy_fingerprint() == stored_output.fingerprint()) .ok_or_else(|| invalid("native maintenance config is absent"))?; + // Raw ingest programs are priced by their stored state components. + if config.derived_input.is_none() { + continue; + } add( format!("maintenance:{}", stored_output.0), json!({"physical_program":program,"stored_output":stored_output,"window_ms":config.stored_window_ms(),"interval_ms":config.slide_interval.saturating_mul(1000)}), diff --git a/control_plane/tests/native_rate_topk.rs b/control_plane/tests/native_rate_topk.rs index 31f22628f..1bbd8b6d4 100644 --- a/control_plane/tests/native_rate_topk.rs +++ b/control_plane/tests/native_rate_topk.rs @@ -55,12 +55,11 @@ fn rate_heap_candidates_bind_durable_counter_windows() { let Ok(plan) = DeploymentPlanCompiler.compile_promql(candidate, environment.clone()) else { continue; }; - if plan - .precompute_plan - .executable_dags - .values() - .any(|dag| !dag.native_programs.is_empty()) - { + if plan.precompute_plan.executable_dags.values().any(|dag| { + dag.native_programs + .keys() + .any(|sink| !dag.reads_raw_samples(*sink)) + }) { continue; } let entry = plan.query_plan.entries.values().next().unwrap(); @@ -151,11 +150,11 @@ fn rate_heap_costs_can_select_each_compiled_candidate() { .map(ToString::to_string) .unwrap_or_default(); let preferred_plan = entry.physical_vector_binding().is_some() - && !plan - .precompute_plan - .executable_dags - .values() - .any(|dag| !dag.native_programs.is_empty()) + && !plan.precompute_plan.executable_dags.values().any(|dag| { + dag.native_programs + .keys() + .any(|sink| !dag.reads_raw_samples(*sink)) + }) && if preferred == "exact" { !program.contains("WithHeap") } else { @@ -218,6 +217,9 @@ fn fixed_window_rate_heap_candidates_install_both_physical_graphs() { }; for installed in plan.precompute_plan.executable_dags.values() { for sink in installed.native_programs.keys() { + if installed.reads_raw_samples(*sink) { + continue; + } let program = installed.native_program(*sink).unwrap().unwrap(); let encoded = String::from_utf8(program.encode().unwrap()).unwrap(); for family in ["CmsWithHeap", "CountSketchWithHeap"] { @@ -278,11 +280,11 @@ fn grouped_rate_placement_follows_summary_store_cost() { continue; } entry.recover_vector_physical_dag().unwrap(); - let stored = plan - .precompute_plan - .executable_dags - .values() - .any(|dag| !dag.native_programs.is_empty()); + let stored = plan.precompute_plan.executable_dags.values().any(|dag| { + dag.native_programs + .keys() + .any(|sink| !dag.reads_raw_samples(*sink)) + }); let rebuilt = program.to_string().contains("SummaryBuild"); assert!(!(stored && rebuilt)); // Planner also offers a relational Sum over the readouts, which has diff --git a/crates/asap_types/src/executable_plan.rs b/crates/asap_types/src/executable_plan.rs index d70baac8b..72daeabfe 100644 --- a/crates/asap_types/src/executable_plan.rs +++ b/crates/asap_types/src/executable_plan.rs @@ -202,6 +202,20 @@ impl InstalledPostAsapDag { Ok(()) } + /// Whether `sink`'s Planner program reads raw ingested samples rather + /// than stored states (a raw-ingest output, not a derived maintenance one). + pub fn reads_raw_samples(&self, sink: PostAsapNodeId) -> bool { + let raw = asap_physical_operators::physical_planner::precompute::raw_sample_schema(); + self.native_program(sink) + .ok() + .flatten() + .is_some_and(|program| { + program + .input_contracts() + .all(|(_, contract)| contract.schema == raw) + }) + } + /// Recovery validates the physical producer's typed storage boundaries; /// it never lowers the semantic provenance document again. pub fn native_program( @@ -235,14 +249,39 @@ impl InstalledPostAsapDag { } } let dag = self.document.decode()?; + let mut raw_inputs = 0; for (id, contract) in program.input_contracts() { let id = PostAsapNodeId(u32::try_from(id).map_err(|_| "physical source id overflow")?); + let node = dag.nodes.iter().find(|n| n.id == id); + // A raw sample boundary is fed by ingestion, not a stored state. + let raw = node.is_some_and(|n| { + asap_physical_operators::physical_planner::precompute::boundary_schema(n) + .is_ok_and(|schema| { + schema == contract.schema + && schema + == asap_physical_operators::physical_planner::precompute::raw_sample_schema() + }) + }); + if raw { + if program.roots().contains(&u64::from(id.0)) + || !matches!( + self.binding.node(id), + Some(BackendNodeBinding::MaintenanceInput) + ) + { + return Err( + "native raw sample source differs from its ingestion binding".into(), + ); + } + raw_inputs += 1; + continue; + } if program.roots().contains(&u64::from(id.0)) || !matches!( self.binding.node(id), Some(BackendNodeBinding::Materialization { .. }) ) - || dag.nodes.iter().find(|n| n.id == id).is_none_or(|n| { + || node.is_none_or(|n| { n.output_schema != *contract.schema && asap_physical_operators::physical_planner::precompute::source_schema( &n.output_schema, @@ -258,6 +297,9 @@ impl InstalledPostAsapDag { if program.input_contracts().count() == 0 { return Err("native precompute program has no bound inputs".into()); } + // Raw summaries are stored as their family's population state; their + // semantic schema may also carry source columns that ingestion drops. + let raw_program = raw_inputs == program.input_contracts().count(); for root in program.roots() { let output = program.output_contract(*root).map_err(|e| e.to_string())?; if dag @@ -265,7 +307,11 @@ impl InstalledPostAsapDag { .iter() .find(|n| u64::from(n.id.0) == *root) .is_none_or(|n| { - n.output_schema != *output.schema + let raw_summary = raw_program + && matches!(&n.payload, planner_types::post_asap::PostAsapOperatorPayload::SummaryAgg { family, .. } + if output.schema == asap_physical_operators::physical_planner::precompute::population_schema(family.clone())); + !raw_summary + && n.output_schema != *output.schema && asap_physical_operators::physical_planner::precompute::source_schema( &n.output_schema, ) diff --git a/data_plane/src/precompute_engine/raw_dag.rs b/data_plane/src/precompute_engine/raw_dag.rs index d549f1533..2180b669c 100644 --- a/data_plane/src/precompute_engine/raw_dag.rs +++ b/data_plane/src/precompute_engine/raw_dag.rs @@ -1,6 +1,15 @@ //! Bind raw ingestion to a selected Planner producer and its raw dependency edge. -use crate::storage_engines::types::KeyByLabelValues; +//! The backend supplies one typed sample batch per pane; the Planner-compiled +//! precompute graph owns every update, grouping and item computation. +use crate::storage_engines::types::{AggregateCore, KeyByLabelValues}; +use asap_physical_operators::{ + operators::Operator, + physical_planner::{precompute, CompiledPhysicalDag, Source as PhysicalSource}, + runtime::{Limits, RunContext, Scope}, + values::{Batch, Value}, +}; use asap_summary_state::factory::{create_planner_accumulator, AccumulatorUpdater}; +use asap_types::physical_plan_codec::PhysicalPlanCodec; use asap_types::{executable_plan::BackendNodeBinding, PrecomputeMaterialization}; use planner_types::post_asap::{ EdgeRole, GroupingStrategy, PostAsapNodeId, PostAsapOperatorPayload, SummaryFamilyType, @@ -19,6 +28,10 @@ pub struct RawDagProgram { 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]>, + source: u64, } impl RawDagProgram { @@ -164,7 +177,21 @@ impl RawDagProgram { if !update_matches { return Err("DAG update differs from stored summary identity".into()); } + let compiled = installed + .native_program(node.id)? + .ok_or("raw materialization lacks its Planner precompute graph")?; + let [(source, contract)] = compiled.input_contracts().collect::>()[..] + else { + return Err("raw precompute graph must read one raw sample input".into()); + }; + if contract.schema != precompute::raw_sample_schema() + || source != u64::from(edge.producer.0) + { + return Err("raw precompute graph does not read the bound raw source".into()); + } let program = Self { + source, + program: compiled.encode().map_err(|e| e.to_string())?.into(), node: node.id, family: family.clone(), input: input.clone(), @@ -195,6 +222,83 @@ impl RawDagProgram { selected.ok_or_else(|| "raw materialization has no selected post-ASAP DAG producer".into()) } + /// 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, + pane: (i64, i64), + max_bytes: usize, + ) -> Result>, String> { + use futures::StreamExt; + 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()?; + for (series, time, value) in samples { + self.apply(&mut *updater, series, value, time)?; + } + return Ok(Some(updater.take_accumulator())); + } + let schema = precompute::raw_sample_schema(); + let rows = samples + .into_iter() + .map(|(series, time, value)| { + precompute::raw_sample_row(&series_labels(series), time, value) + }) + .collect(); + let batch = Batch::try_new(schema.clone(), rows).map_err(|e| e.to_string())?; + let sources = std::collections::BTreeMap::from([( + self.source, + Box::new(Operator::source(schema, vec![batch]).map_err(|e| e.to_string())?) + as PhysicalSource<'_>, + )]); + let program = decoded(&self.program)?; + let graph = program.instantiate(sources).map_err(|e| e.to_string())?; + let context = RunContext::new( + Scope::Ingestion { + window_start_ms: pane.0, + window_end_ms: pane.1, + revision: 0, + }, + Limits { + max_bytes, + ..Limits::default() + }, + ) + .map_err(|e| e.to_string())?; + let rows = futures::executor::block_on(async { + let mut stream = graph.execute(program.roots(), context)?.pop().ok_or( + asap_physical_operators::Error::Invalid("missing output".into()), + )?; + let mut rows = Vec::new(); + while let Some(batch) = stream.next().await { + 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) } @@ -289,3 +393,59 @@ impl RawDagProgram { Ok(()) } } + +/// The complete label set of a canonical series key, including `__name__`. +pub(crate) fn series_labels(series: &str) -> std::collections::BTreeMap { + let mut labels: std::collections::BTreeMap<_, _> = + super::worker::parse_labels_from_series_key(series) + .into_iter() + .map(|(k, v)| { + ( + k.to_owned(), + super::worker::decode_label_value(v).into_owned(), + ) + }) + .collect(); + let metric = series.split('{').next().unwrap_or_default(); + if !metric.is_empty() { + labels.insert("__name__".into(), metric.to_owned()); + } + labels +} + +/// A pane's samples as ingestion delivers them to Planner: in timestamp +/// order (stable for equal times), with one sample per series and timestamp; +/// a repeated `(series, timestamp)` is the same sample, so the first is kept. +fn pane_batch<'a>( + samples: impl IntoIterator, +) -> Vec<(&'a str, i64, f64)> { + let mut samples = samples.into_iter().collect::>(); + samples.sort_by_key(|(_, time, _)| *time); + let mut seen = std::collections::HashSet::new(); + samples.retain(|(series, time, _)| seen.insert((*series, *time))); + samples +} + +/// Decoded graphs are not `Send`, so each thread decodes an installed graph +/// once. Holding the encoded bytes keeps their address from being reused. +fn decoded(encoded: &std::sync::Arc<[u8]>) -> Result, String> { + type Cache = + std::collections::HashMap, std::rc::Rc)>; + thread_local! { + static DECODED: std::cell::RefCell = Default::default(); + } + DECODED.with(|cache| { + let mut cache = cache.borrow_mut(); + let key = encoded.as_ptr() as usize; + if let Some((_, program)) = cache.get(&key) { + return Ok(program.clone()); + } + if cache.len() >= 1024 { + cache.clear(); + } + let program = + std::rc::Rc::new(CompiledPhysicalDag::decode(encoded).map_err(|e| e.to_string())?); + cache.insert(key, (encoded.clone(), program.clone())); + Ok(program) + }) +} diff --git a/data_plane/src/precompute_engine/worker.rs b/data_plane/src/precompute_engine/worker.rs index 2f6d3e085..814ed12d0 100644 --- a/data_plane/src/precompute_engine/worker.rs +++ b/data_plane/src/precompute_engine/worker.rs @@ -57,8 +57,9 @@ struct GroupState { /// `group_key` field. group_key: Arc, window_manager: WindowManager, - /// Active panes for raw-sample accumulation, keyed by pane_start_ms. - active_panes: BTreeMap>, + /// Samples admitted to each open pane, keyed by pane_start_ms, in arrival + /// order. The installed Planner graph builds the pane's state at close. + active_panes: BTreeMap, i64, f64)>>, /// Last cumulative counter sample per source series for heap membership /// materializations. This is bounded O(series) derivative state, not a /// raw-sample history, and deliberately survives pane rotation. @@ -629,6 +630,8 @@ impl Worker { // layouts update one non-overlapping base pane; FullWindow updates // every overlapping semantic window that contains the sample. for (series_key, ts, val) in &samples { + // Shared by every pane (overlapping full windows) the sample enters. + let series: Arc = Arc::from(series_key.as_str()); let too_late = previous_event_time != i64::MIN && pane_timestamp(*ts) < watermark_for_event_time(previous_event_time, allowed_lateness_ms); @@ -691,16 +694,15 @@ impl Worker { continue; } record_late_input("append_correction", "raw_sample"); - let mut updater = - installed_updater(state.program.as_deref(), &state.config)?; - apply_installed_sample( + let Some(correction) = build_pane( state.program.as_deref(), - &mut *updater, - series_key, - *val, - *ts, &state.config, - )?; + &[(Arc::from(series_key.as_str()), *ts, *val)], + (bucket_start, bucket_end), + )? + else { + continue; + }; if let (Some(observer), Some(revision)) = (&self.erp_observer, &input_revision) { @@ -731,7 +733,7 @@ impl Worker { state.catalog_generation.as_ref(), state.stored_output_reference.clone(), ); - emit_batch.push((output, updater.take_accumulator())); + emit_batch.push((output, correction)); debug!( "Forwarding late sample to store for evicted pane [{}, {})", bucket_start, bucket_end @@ -746,21 +748,9 @@ impl Worker { // only closes an idle pane, not a long-running bulk ingest whose // records share one event timestamp. state.touch_pane(bucket_start, now_ms); - if let std::collections::btree_map::Entry::Vacant(entry) = - state.active_panes.entry(bucket_start) - { - entry.insert(installed_updater(state.program.as_deref(), &state.config)?); - } - let updater = state.active_panes.get_mut(&bucket_start).unwrap(); + let pane = state.active_panes.entry(bucket_start).or_default(); if let Some(value) = value { - apply_installed_sample( - state.program.as_deref(), - &mut **updater, - series_key, - value, - *ts, - &state.config, - )?; + pane.push((Arc::clone(&series), *ts, value)); if let (Some(observer), Some(revision)) = (&self.erp_observer, &input_revision) { observer.observe( @@ -792,10 +782,8 @@ impl Worker { for window_start in &closed { let (_, window_end) = state.bucket_bounds(*window_start); - let pane_starts = [*window_start]; - if let Some(accumulator) = merge_panes_for_window(&mut state.active_panes, &pane_starts) - { + if let Some(accumulator) = close_pane(state, *window_start)? { let key = build_group_key_label_values(group_key); let output = precomputed_output_for_group( *window_start as u64, @@ -971,12 +959,10 @@ impl Worker { let closed = state.closed_buckets(previous_closure_watermark, event_watermark); for window_start in &closed { let (_, window_end) = state.bucket_bounds(*window_start); - let pane_starts = [*window_start]; // Emit from the raw-sample pane map (in case both sources are // populated for the same group; rare but supported). - if let Some(accumulator) = merge_panes_for_window(&mut state.active_panes, &pane_starts) - { + if let Some(accumulator) = close_pane(state, *window_start)? { let key = build_group_key_label_values(group_key); let output = precomputed_output_for_group( *window_start as u64, @@ -993,9 +979,7 @@ impl Worker { } // Emit from the sketch pane map. - if let Some(accumulator) = - merge_sketch_panes_for_window(&mut state.sketch_panes, &pane_starts) - { + if let Some(accumulator) = state.sketch_panes.remove(window_start) { let key = build_group_key_label_values(group_key); let output = precomputed_output_for_group( *window_start as u64, @@ -1177,11 +1161,8 @@ impl Worker { for window_start in &closed { let (_, window_end) = state.bucket_bounds(*window_start); - let pane_starts = [*window_start]; - if let Some(accumulator) = - merge_panes_for_window(&mut state.active_panes, &pane_starts) - { + if let Some(accumulator) = close_pane(state, *window_start)? { let key = build_group_key_label_values(&group_key); let output = precomputed_output_for_group( *window_start as u64, @@ -1197,9 +1178,7 @@ impl Worker { emit_batch.push((output, accumulator)); } - if let Some(accumulator) = - merge_sketch_panes_for_window(&mut state.sketch_panes, &pane_starts) - { + if let Some(accumulator) = state.sketch_panes.remove(window_start) { let key = build_group_key_label_values(&group_key); let output = precomputed_output_for_group( *window_start as u64, @@ -1285,11 +1264,8 @@ impl Worker { for window_start in &closed { let (_, window_end) = state.bucket_bounds(*window_start); - let pane_starts = [*window_start]; - if let Some(accumulator) = - merge_panes_for_window(&mut state.active_panes, &pane_starts) - { + if let Some(accumulator) = close_pane(state, *window_start)? { let key = build_group_key_label_values(&group_key); let output = precomputed_output_for_group( *window_start as u64, @@ -1305,9 +1281,7 @@ impl Worker { emit_batch.push((output, accumulator)); } - if let Some(accumulator) = - merge_sketch_panes_for_window(&mut state.sketch_panes, &pane_starts) - { + if let Some(accumulator) = state.sketch_panes.remove(window_start) { let key = build_group_key_label_values(&group_key); let output = precomputed_output_for_group( *window_start as u64, @@ -1614,45 +1588,54 @@ pub fn decode_label_value(s: &str) -> std::borrow::Cow<'_, str> { std::borrow::Cow::Owned(out) } -fn installed_updater( +/// Build a pane's state from its samples with the installed Planner graph. +fn build_pane( program: Option<&super::raw_dag::RawDagProgram>, config: &PrecomputeMaterialization, -) -> Result, String> { + samples: &[(Arc, i64, f64)], + pane: (i64, i64), +) -> Result>, String> { if let Some(program) = program { - return program.updater(); + return program.build( + samples + .iter() + .map(|(series, time, value)| (series.as_ref(), *time, *value)), + pane, + asap_physical_operators::runtime::Limits::default().max_bytes, + ); } #[cfg(test)] { - Ok(create_fixture_accumulator(config)) + let _ = pane; + let mut updater = create_fixture_accumulator(config); + for (series, time, value) in samples { + apply_sample(&mut *updater, series, *value, *time, config); + } + Ok(Some(updater.take_accumulator())) } #[cfg(not(test))] { - let _ = config; + let _ = (config, samples, pane); Err("missing installed Planner producer".into()) } } -fn apply_installed_sample( - program: Option<&super::raw_dag::RawDagProgram>, - updater: &mut dyn AccumulatorUpdater, - series: &str, - value: f64, - timestamp: i64, - config: &PrecomputeMaterialization, -) -> Result<(), String> { - if let Some(program) = program { - return program.apply(updater, series, value, timestamp); - } - #[cfg(test)] - { - apply_sample(updater, series, value, timestamp, config); - Ok(()) - } - #[cfg(not(test))] - { - let _ = config; - Err("missing installed Planner producer".into()) - } +/// Build a closed pane's state from every sample it admitted, then remove it. +fn close_pane( + state: &mut GroupState, + start: i64, +) -> Result>, String> { + let Some(samples) = state.active_panes.get(&start) else { + return Ok(None); + }; + let built = build_pane( + state.program.as_deref(), + &state.config, + samples, + state.bucket_bounds(start), + )?; + state.active_panes.remove(&start); + Ok(built) } /// Route a single sample to `updater`, dispatching keyed vs. non-keyed based on config. @@ -1755,75 +1738,6 @@ fn extract_aggregated_key_from_series( KeyByLabelValues::new_with_labels(values) } -/// Merge the pane accumulators that constitute a closed window. -/// -/// The oldest pane (index 0) is taken destructively from `active_panes` -/// (no future window needs it). All later panes are snapshot-read -/// (non-destructive; they are shared by newer overlapping windows). -/// -/// Returns `None` if all panes for the window are absent. -fn merge_panes_for_window( - active_panes: &mut BTreeMap>, - pane_starts: &[i64], -) -> Option> { - let mut merged: Option> = None; - - for (i, &ps) in pane_starts.iter().enumerate() { - let pane_acc = if i == 0 { - // Oldest pane: evict and MOVE the accumulator out (no clone). - active_panes - .remove(&ps) - .map(|updater| updater.into_accumulator()) - } else { - // Shared pane: non-destructive snapshot - active_panes - .get(&ps) - .map(|updater| updater.snapshot_accumulator()) - }; - - if let Some(acc) = pane_acc { - merged = Some(match merged { - None => acc, - Some(existing) => existing.merge_with(acc.as_ref()).unwrap_or(existing), - }); - } - } - - merged -} - -/// Merge pre-built accumulator panes for a window. -/// -/// Equivalent to `merge_panes_for_window` but operating on the sketch pane -/// map (`Box` directly). The oldest pane is destructively -/// taken (it will never be needed by a later window); subsequent panes are -/// cloned so that still-open overlapping windows can still read them. -fn merge_sketch_panes_for_window( - sketch_panes: &mut BTreeMap>, - pane_starts: &[i64], -) -> Option> { - let mut merged: Option> = None; - - for (i, &ps) in pane_starts.iter().enumerate() { - let pane_acc: Option> = if i == 0 { - // Oldest pane: destructive take + evict - sketch_panes.remove(&ps) - } else { - // Shared pane: non-destructive clone - sketch_panes.get(&ps).map(|acc| acc.clone_boxed_core()) - }; - - if let Some(acc) = pane_acc { - merged = Some(match merged { - None => acc, - Some(existing) => existing.merge_with(acc.as_ref()).unwrap_or(existing), - }); - } - } - - merged -} - #[cfg(test)] mod tests { use super::*; @@ -4685,6 +4599,242 @@ mod dag_execution_tests { } } + fn plan_with( + query: &str, + accuracy: Option, + entry: usize, + ) -> Option { + let mut json: serde_json::Value = serde_json::from_str(include_str!( + "../../../docs/examples/asapquery-compatibility-demo-snapshot.json" + )) + .unwrap(); + let mut item = json["query_workload"]["repeating_queries"][entry].clone(); + item["query"] = query.into(); + if let Some(accuracy) = accuracy { + item["requirements"]["accuracy"] = accuracy; + } + json["query_workload"]["repeating_queries"] = serde_json::json!([item]); + let snapshot = serde_json::from_value(json).ok()?; + crate::tests::test_utilities::planning::quoted_snapshot(snapshot, false) + .compile_promql() + .ok() + } + + /// Queries, accuracies and snapshot entries whose plans install raw outputs + /// of every exact family, DDSketch and grouped variants. + fn raw_output_fixtures() -> Vec { + let accuracies = [ + None, + Some(serde_json::json!({"explicit": {"Epsilon": 0.05}})), + Some( + serde_json::json!({"explicit": {"EpsilonDelta": {"epsilon": 0.01, "delta": 0.01}}}), + ), + ]; + let mut plans = Vec::new(); + for query in [ + "sum_over_time(asap_demo_gauge[5s])", + "count_over_time(asap_demo_gauge[5s])", + "min_over_time(asap_demo_gauge[5s])", + "max_over_time(asap_demo_gauge[5s])", + "rate(asap_demo_counter_total[5s])", + "increase(asap_demo_counter_total[5s])", + "sum by (service) (sum_over_time(asap_demo_gauge[5s]))", + "sum by (service) (rate(asap_demo_counter_total[5s]))", + "quantile_over_time(0.99, asap_demo_gauge[5s])", + "sum by (service) (quantile_over_time(0.99, asap_demo_gauge[5s]))", + "quantile by (service) (0.9, asap_demo_gauge)", + "distinct_over_time(asap_demo_gauge[5s])", + "count(asap_demo_gauge)", + "topk(2, sum_over_time(asap_demo_gauge[5s]))", + "topk(2, rate(asap_demo_counter_total[5s]))", + "topk(2, asap_demo_gauge)", + ] { + for accuracy in &accuracies { + for entry in [0, 3] { + plans.extend(plan_with(query, accuracy.clone(), entry)); + } + } + } + plans + } + + // Live raw ingest builds each pane with the installed Planner graph, and + // every stored state equals feeding the selected kernel sample by sample. + #[test] + fn live_panes_execute_planner_dag_with_identical_states() { + let mut outputs = 0; + let mut families = std::collections::BTreeSet::new(); + for plan in raw_output_fixtures() { + let installed = + InstalledPrecomputePlan::from_precompute_plan(plan.precompute_plan.clone()) + .unwrap(); + for config in plan + .precompute_plan + .materializations + .iter() + .filter(|c| c.derived_input.is_none()) + { + let fp = config.policy_fingerprint(); + let program = installed.raw_programs[&fp.as_u64()].clone(); + let series_scoped = config.partitioning + == Some(asap_types::sds::PopulationPartitioning::PerEntity) + || (config.partitioning.is_none() + && matches!( + config.aggregation_type, + asap_types::AggregationType::Increase + | asap_types::AggregationType::Rate + | asap_types::AggregationType::Min + | asap_types::AggregationType::Max + )); + // Routing identity as Remote Write assigns it. + let mut groups = BTreeMap::, Vec<(String, i64, f64)>>::new(); + for time in (500..=9500).step_by(1000) { + for (index, (service, instance)) in + [("a", "1"), ("a", "2"), ("b", "1")].into_iter().enumerate() + { + let labels = BTreeMap::from([ + ("instance".to_string(), instance.to_string()), + ("service".to_string(), service.to_string()), + ]); + let key = if series_scoped { + labels.clone().into_iter().collect() + } else { + config + .grouping_labels + .iter() + .map(|n| (n.clone(), labels.get(n).cloned().unwrap_or_default())) + .collect() + }; + let value = (index as f64 + 1.0) * time as f64 / 100.0; + groups.entry(key).or_default().push(( + format!( + "{}{{instance=\"{instance}\",service=\"{service}\"}}", + config.metric + ), + time, + value, + )); + } + } + let sink = Arc::new(CapturingOutputSink::new()); + let (_tx, rx) = mpsc::channel(8); + let mut worker = Worker::new( + 0, + rx, + sink.clone(), + InstalledPrecomputePlanHandle::new(installed.clone()), + WorkerRuntimeConfig { + max_buffer_per_series: 100, + allowed_lateness_ms: 0, + pass_raw_samples: false, + raw_mode_aggregation_id: 0, + late_data_policy: LateDataPolicy::Drop, + wall_clock_idle_grace_period_ms: 0, + wall_clock_max_open_grace_period_ms: 0, + }, + Arc::new(AtomicUsize::new(0)), + Arc::new(AtomicI64::new(0)), + ); + let manager = WindowManager::with_layout( + config.window_size, + config.slide_interval, + config.pane_origin_ms, + &config.window_layout, + ); + let right_closed = config + .parameters + .get("promql_right_closed") + .and_then(serde_json::Value::as_bool) + .unwrap_or(false); + let mut expected = 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())), + )); + for (series, time, value) in samples { + let pane_time = if right_closed { time - 1 } else { *time }; + for start in manager.stored_bucket_starts(pane_time) { + let (start, end) = manager.stored_bucket_bounds(start); + let updater = expected + .entry(( + group_key.values().labels.join(";"), + start as u64, + end as u64, + )) + .or_insert_with(|| program.updater().unwrap()); + program + .apply(&mut **updater, series, *value, *time) + .unwrap(); + } + } + } + // One batch per scrape, so the watermark closes earlier panes + // while later ones are open; shutdown closes the rest. + for time in (500..=9500).step_by(1000) { + for (sid, (key, samples)) in groups.iter().enumerate() { + let group_key = Arc::new(GroupKey::new( + key.iter().map(|(k, v)| (k.as_str(), v.as_str())), + )); + let batch = samples.iter().filter(|s| s.1 == time).cloned().collect(); + worker + .process_group_samples(sid as u64 + 1, fp, &group_key, batch) + .unwrap(); + } + } + worker.force_close_all().unwrap(); + let actual = sink + .drain() + .into_iter() + .map(|(output, state)| { + let key = output + .key + .as_ref() + .map(|k| k.labels.join(";")) + .unwrap_or_default(); + ( + (key, output.start_timestamp, output.end_timestamp), + state.serialize_to_bytes(), + ) + }) + .collect::>(); + let expected = expected + .into_iter() + .map(|(key, mut updater)| { + (key, updater.take_accumulator().serialize_to_bytes()) + }) + .collect::>(); + assert_eq!( + actual.keys().collect::>(), + expected.keys().collect::>(), + "{}", + config.metric + ); + assert_eq!(actual, expected, "{:?}", config.aggregation_type); + outputs += 1; + families.insert( + format!("{:?}", config.accumulator_spec().unwrap().family) + .split(['(', ' ', ',']) + .find(|part| { + [ + "Sum", "Count", "Min", "Max", "Rate", "Increase", "DDSketch", + "Kll", "Hll", + ] + .contains(part) + }) + .unwrap_or("other") + .to_owned(), + ); + } + } + assert!(outputs >= 50, "only {outputs} raw outputs exercised"); + for family in ["Sum", "Count", "Min", "Max", "Rate", "Increase", "DDSketch"] { + assert!( + families.contains(family), + "no {family} output: {families:?}" + ); + } + } + // A flat config and a DAG whose producer no longer matches its binding cannot install. #[test] fn execution_requires_matching_dag_producer() { @@ -4701,6 +4851,153 @@ mod dag_execution_tests { .to_string() .contains("DAG producer")); } + fn single_worker( + plan: &control_plane::physical::compiler::CompiledPhysicalPlan, + sink: Arc, + late_data_policy: LateDataPolicy, + ) -> Worker { + let (_tx, rx) = mpsc::channel(8); + Worker::new( + 0, + rx, + sink, + InstalledPrecomputePlanHandle::new( + InstalledPrecomputePlan::from_precompute_plan(plan.precompute_plan.clone()) + .unwrap(), + ), + WorkerRuntimeConfig { + max_buffer_per_series: 100, + allowed_lateness_ms: 0, + pass_raw_samples: false, + raw_mode_aggregation_id: 0, + late_data_policy, + wall_clock_idle_grace_period_ms: 0, + wall_clock_max_open_grace_period_ms: 0, + }, + Arc::new(AtomicUsize::new(0)), + Arc::new(AtomicI64::new(0)), + ) + } + + // A counter pane is delivered to Planner in timestamp order with one sample + // per series and timestamp; arrival order and a resent sample do not fail it. + #[test] + fn counter_pane_orders_and_deduplicates_samples() { + let plan = plan("rate(asap_demo_counter_total[5s])"); + let config = plan.precompute_plan.materializations[0].clone(); + let program = InstalledPrecomputePlan::from_precompute_plan(plan.precompute_plan.clone()) + .unwrap() + .raw_programs[&config.policy_fp_u64()] + .clone(); + let series = format!("{}{{job=\"api\"}}", config.metric); + let arrived = [ + (1100, 10.0), + (1300, 30.0), + (1200, 20.0), + (1300, 99.0), + (1400, 40.0), + ]; + let sink = Arc::new(CapturingOutputSink::new()); + let mut worker = single_worker(&plan, sink.clone(), LateDataPolicy::Drop); + worker + .process_group_samples( + 1, + config.policy_fingerprint(), + &Arc::new(GroupKey::new([("job", "api")])), + arrived + .iter() + .map(|(t, v)| (series.clone(), *t, *v)) + .collect(), + ) + .unwrap(); + worker.force_close_all().unwrap(); + let manager = WindowManager::with_layout( + config.window_size, + config.slide_interval, + config.pane_origin_ms, + &config.window_layout, + ); + let right_closed = config + .parameters + .get("promql_right_closed") + .and_then(serde_json::Value::as_bool) + .unwrap_or(false); + let mut panes = BTreeMap::<(i64, i64), Vec<(&str, i64, f64)>>::new(); + for (time, value) in [(1100, 10.0), (1200, 20.0), (1300, 30.0), (1400, 40.0)] { + let pane_time = if right_closed { time - 1 } else { time }; + for start in manager.stored_bucket_starts(pane_time) { + panes + .entry(manager.stored_bucket_bounds(start)) + .or_default() + .push((series.as_str(), time, value)); + } + } + let actual = sink + .drain() + .into_iter() + .map(|(output, state)| { + ( + (output.start_timestamp as i64, output.end_timestamp as i64), + state.serialize_to_bytes(), + ) + }) + .collect::>(); + let expected = panes + .into_iter() + .map(|(pane, samples)| { + let state = program.build(samples, pane, 1 << 20).unwrap().unwrap(); + (pane, state.serialize_to_bytes()) + }) + .collect::>(); + assert!(!expected.is_empty()); + assert_eq!(actual, expected); + } + + // A late sample forwarded to a closed pane is its own Planner-built correction. + #[test] + fn late_forwarded_sample_is_built_by_the_planner_graph() { + let plan = plan("sum_over_time(asap_demo_gauge[5s])"); + let config = plan.precompute_plan.materializations[0].clone(); + let sink = Arc::new(CapturingOutputSink::new()); + let mut worker = single_worker(&plan, sink.clone(), LateDataPolicy::ForwardToStore); + let group = Arc::new(GroupKey::new([])); + let fp = config.policy_fingerprint(); + let series = config.metric.clone(); + for time in [1000, 12_000] { + worker + .process_group_samples(1, fp, &group, vec![(series.clone(), time, 2.0)]) + .unwrap(); + } + sink.drain(); + worker + .process_group_samples(1, fp, &group, vec![(series.clone(), 1500, 7.0)]) + .unwrap(); + let corrections = sink.drain(); + assert!(!corrections.is_empty(), "late sample must be forwarded"); + for (_, state) in corrections { + let exact = + ExactAccumulator::deserialize_from_bytes(&state.serialize_to_bytes()).unwrap(); + assert_eq!( + exact + .query_statistic(asap_types::Statistic::Sum, &None, &Default::default()) + .unwrap(), + 7.0 + ); + } + } + + // A raw output installs only with its Planner-compiled precompute graph. + #[test] + fn raw_output_requires_its_planner_precompute_graph() { + let mut plan = plan("sum_over_time(asap_demo_gauge[5s])").precompute_plan; + for installed in plan.executable_dags.values_mut() { + installed.native_programs.clear(); + } + assert!(InstalledPrecomputePlan::from_precompute_plan(plan) + .unwrap_err() + .to_string() + .contains("Planner precompute graph")); + } // Changing a raw update must not silently reuse the original summary identity. #[test] fn altered_dag_update_cannot_reuse_a_stored_definition() { diff --git a/data_plane/tests/support/native_revision_process.rs b/data_plane/tests/support/native_revision_process.rs index 32e7eab0e..2580fd490 100644 --- a/data_plane/tests/support/native_revision_process.rs +++ b/data_plane/tests/support/native_revision_process.rs @@ -76,10 +76,12 @@ async fn run_native_ensemble(family: Option<&str>, cadence_ms: u64, phase_ms: u6 .into_physical_compilation_request() .unwrap(); let preferred = |plan: &control_plane::physical::compiler::CompiledPhysicalPlan| { + // Raw outputs also carry Planner programs; the ensemble needs a derived one. plan.precompute_plan.executable_dags.values().any(|dag| { - dag.native_programs - .values() - .any(|program| family.is_none_or(|family| program.to_string().contains(family))) + dag.native_programs.iter().any(|(sink, program)| { + !dag.reads_raw_samples(*sink) + && family.is_none_or(|family| program.to_string().contains(family)) + }) }) && plan.query_plan.entries.values().all(|entry| { !entry.nodes.values().any(|node| { matches!(