From 07489ba1e3e6143b7e1e58a38cff04aa64c0b8f4 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 07:15:11 +0000 Subject: [PATCH 1/4] build: repin Planner to the integration revision with raw-sample precompute Co-Authored-By: Claude Opus 5.5 --- Cargo.lock | 12 ++++++------ Cargo.toml | 11 ++++++----- 2 files changed, 12 insertions(+), 11 deletions(-) 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"] } + From 87e9c16a2ec78aa4f5e94c0f87a47fb13e03acd6 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 07:15:11 +0000 Subject: [PATCH 2/4] feat(control-plane): retain Planner precompute graphs for raw outputs Each raw stored output's native program is Planner's precompute graph from its raw sample boundary to the output. Validation admits raw sample inputs bound as ingestion inputs; pricing keeps charging raw outputs by their stored state. Tests that meant "a derived maintenance program exists" now skip raw-ingest programs. Co-Authored-By: Claude Opus 5.5 --- .../src/physical/executable_binding.rs | 20 ++++++++ control_plane/src/physical/workload_cost.rs | 4 ++ control_plane/tests/native_rate_topk.rs | 34 +++++++------ crates/asap_types/src/executable_plan.rs | 50 ++++++++++++++++++- .../tests/support/native_revision_process.rs | 8 +-- 5 files changed, 95 insertions(+), 21 deletions(-) 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/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/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!( From d1723e64adac9d7b32dbb342a33af39564bd0038 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 07:15:11 +0000 Subject: [PATCH 3/4] refactor(precompute): build live raw panes with the Planner graph The worker keeps each open pane's admitted samples and, when the pane closes (or a late sample is forwarded), executes the output's installed Planner graph over them as one typed batch. Pane assignment, completeness and lateness stay in the worker. The single-pane merge helpers become a plain remove. Heaps keep their kernel until stored heap readout decodes Planner's weighted frequency state. Co-Authored-By: Claude Opus 5.5 --- data_plane/src/precompute_engine/raw_dag.rs | 124 +++++- data_plane/src/precompute_engine/worker.rs | 434 +++++++++++++------- 2 files changed, 414 insertions(+), 144 deletions(-) diff --git a/data_plane/src/precompute_engine/raw_dag.rs b/data_plane/src/precompute_engine/raw_dag.rs index d549f1533..6bfde2ce2 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,82 @@ 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 in arrival order. + /// 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; + // 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 = CompiledPhysicalDag::decode(&self.program).map_err(|e| e.to_string())?; + 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 +392,22 @@ 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 +} diff --git a/data_plane/src/precompute_engine/worker.rs b/data_plane/src/precompute_engine/worker.rs index 2f6d3e085..f887e932a 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>, /// 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. @@ -691,16 +692,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, - )?; + &[(series_key.clone(), *ts, *val)], + (bucket_start, bucket_end), + )? + else { + continue; + }; if let (Some(observer), Some(revision)) = (&self.erp_observer, &input_revision) { @@ -731,7 +731,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 +746,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((series_key.clone(), *ts, value)); if let (Some(observer), Some(revision)) = (&self.erp_observer, &input_revision) { observer.observe( @@ -792,10 +780,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 +957,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 +977,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 +1159,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 +1176,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 +1262,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 +1279,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 +1586,52 @@ 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: &[(String, 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_str(), *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()) - } +/// Remove a closed pane and build its state from every sample it admitted. +fn close_pane( + state: &mut GroupState, + start: i64, +) -> Result>, String> { + let Some(samples) = state.active_panes.remove(&start) else { + return Ok(None); + }; + build_pane( + state.program.as_deref(), + &state.config, + &samples, + state.bucket_bounds(start), + ) } /// Route a single sample to `updater`, dispatching keyed vs. non-keyed based on config. @@ -1755,75 +1734,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 +4595,232 @@ 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: 60_000, + 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 (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())), + )); + 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(); + } + } + worker + .process_group_samples(sid as u64 + 1, fp, &group_key, samples.clone()) + .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 +4837,18 @@ mod dag_execution_tests { .to_string() .contains("DAG producer")); } + // 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() { From 2edfc571860330289026d72cd912a8497dc73095 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 07:31:09 +0000 Subject: [PATCH 4/4] fix(precompute): deliver time-ordered, deduplicated pane batches A pane's samples reach the Planner graph in timestamp order with one sample per series and timestamp, so out-of-order arrival or a resent sample no longer fails a counter pane. A pane is removed only after it builds; each thread decodes an installed graph once; buffered samples share their series key. Tests cover watermark closure across scrape batches, counter ordering and late forwarded corrections. Co-Authored-By: Claude Opus 5.5 --- control_plane/src/physical/plan_dot.rs | 6 +- data_plane/src/precompute_engine/raw_dag.rs | 42 ++++- data_plane/src/precompute_engine/worker.rs | 179 ++++++++++++++++++-- 3 files changed, 209 insertions(+), 18 deletions(-) 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/data_plane/src/precompute_engine/raw_dag.rs b/data_plane/src/precompute_engine/raw_dag.rs index 6bfde2ce2..2180b669c 100644 --- a/data_plane/src/precompute_engine/raw_dag.rs +++ b/data_plane/src/precompute_engine/raw_dag.rs @@ -222,7 +222,7 @@ 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 in arrival order. + /// 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, @@ -231,6 +231,7 @@ impl RawDagProgram { 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!( @@ -257,7 +258,7 @@ impl RawDagProgram { Box::new(Operator::source(schema, vec![batch]).map_err(|e| e.to_string())?) as PhysicalSource<'_>, )]); - let program = CompiledPhysicalDag::decode(&self.program).map_err(|e| e.to_string())?; + let program = decoded(&self.program)?; let graph = program.instantiate(sources).map_err(|e| e.to_string())?; let context = RunContext::new( Scope::Ingestion { @@ -411,3 +412,40 @@ pub(crate) fn series_labels(series: &str) -> std::collections::BTreeMap( + 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 f887e932a..814ed12d0 100644 --- a/data_plane/src/precompute_engine/worker.rs +++ b/data_plane/src/precompute_engine/worker.rs @@ -59,7 +59,7 @@ struct GroupState { window_manager: WindowManager, /// 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>, + 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. @@ -630,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); @@ -695,7 +697,7 @@ impl Worker { let Some(correction) = build_pane( state.program.as_deref(), &state.config, - &[(series_key.clone(), *ts, *val)], + &[(Arc::from(series_key.as_str()), *ts, *val)], (bucket_start, bucket_end), )? else { @@ -748,7 +750,7 @@ impl Worker { state.touch_pane(bucket_start, now_ms); let pane = state.active_panes.entry(bucket_start).or_default(); if let Some(value) = value { - pane.push((series_key.clone(), *ts, value)); + pane.push((Arc::clone(&series), *ts, value)); if let (Some(observer), Some(revision)) = (&self.erp_observer, &input_revision) { observer.observe( @@ -1590,14 +1592,14 @@ pub fn decode_label_value(s: &str) -> std::borrow::Cow<'_, str> { fn build_pane( program: Option<&super::raw_dag::RawDagProgram>, config: &PrecomputeMaterialization, - samples: &[(String, i64, f64)], + samples: &[(Arc, i64, f64)], pane: (i64, i64), ) -> Result>, String> { if let Some(program) = program { return program.build( samples .iter() - .map(|(series, time, value)| (series.as_str(), *time, *value)), + .map(|(series, time, value)| (series.as_ref(), *time, *value)), pane, asap_physical_operators::runtime::Limits::default().max_bytes, ); @@ -1618,20 +1620,22 @@ fn build_pane( } } -/// Remove a closed pane and build its state from every sample it admitted. +/// 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.remove(&start) else { + let Some(samples) = state.active_panes.get(&start) else { return Ok(None); }; - build_pane( + let built = build_pane( state.program.as_deref(), &state.config, - &samples, + 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. @@ -4721,7 +4725,7 @@ mod dag_execution_tests { InstalledPrecomputePlanHandle::new(installed.clone()), WorkerRuntimeConfig { max_buffer_per_series: 100, - allowed_lateness_ms: 60_000, + allowed_lateness_ms: 0, pass_raw_samples: false, raw_mode_aggregation_id: 0, late_data_policy: LateDataPolicy::Drop, @@ -4743,7 +4747,7 @@ mod dag_execution_tests { .and_then(serde_json::Value::as_bool) .unwrap_or(false); let mut expected = BTreeMap::new(); - for (sid, (key, samples)) in groups.iter().enumerate() { + for (key, samples) in &groups { let group_key = Arc::new(GroupKey::new( key.iter().map(|(k, v)| (k.as_str(), v.as_str())), )); @@ -4763,9 +4767,19 @@ mod dag_execution_tests { .unwrap(); } } - worker - .process_group_samples(sid as u64 + 1, fp, &group_key, samples.clone()) - .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 @@ -4837,6 +4851,141 @@ 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() {