From ee9e0b15cafb9c39dfb3bc724eaa1ee3ed8da39b Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 02:22:53 +0000 Subject: [PATCH 1/2] refactor!: move pane construction out of the physical layer Pane geometry, pane population checks and per-pane scheduling are deployment concerns. The physical layer keeps the computation the deployment binds: per-input summary build, union, shared merge and readout. - Remove `physical_planner::compile_temporal_pane_candidate` and its `TemporalPaneMaintenance`/`TemporalPaneCandidate`/`TemporalEntityIdentity` contract. - Remove the `PaneInput` operator; `ScopeTimestamp` remains as a general operator in `operators/scope_timestamp`. - Drop the `selected_temporal_lifecycle_compiles_panes_and_executes` E2E test and update the crate README and design doc. Co-Authored-By: Claude Opus 5.5 --- crates/asap-physical-operators/README.md | 15 +- .../src/operators/mod.rs | 18 +- .../src/operators/panes.rs | 228 --------- .../src/operators/scope_timestamp.rs | 91 ++++ .../src/operators/unchecked.rs | 5 - .../src/physical_planner/mod.rs | 6 - .../src/physical_planner/temporal_panes.rs | 331 ------------- .../summary_maintenance_lifecycle_e2e.rs | 443 ------------------ .../physical-planning-and-deployment.md | 27 +- 9 files changed, 103 insertions(+), 1061 deletions(-) delete mode 100644 crates/asap-physical-operators/src/operators/panes.rs create mode 100644 crates/asap-physical-operators/src/operators/scope_timestamp.rs delete mode 100644 crates/asap-physical-operators/src/physical_planner/temporal_panes.rs diff --git a/crates/asap-physical-operators/README.md b/crates/asap-physical-operators/README.md index 5111d316..5f638465 100644 --- a/crates/asap-physical-operators/README.md +++ b/crates/asap-physical-operators/README.md @@ -116,15 +116,6 @@ the shared runtime with independent per-run state. Window coverage, revision and maintenance-policy admission remain deployment/planning contracts; this compiler does not discover storage or silently change a selected maintenance strategy. -`physical_planner::compile_temporal_pane_candidate` lowers a selected continuous -KLL lifecycle and Sliding/Tumbling framework into maintenance and query DAGs. -`TemporalPaneMaintenance` supplies pane geometry and a resolved complete entity -identity contract. The compiler inserts population guards, scan predicates, -pane construction, ordered state slots, a shared merge and quantile readouts. -Pane outputs have distinct physical identities from the logical whole-window -summary, and the returned candidate retains the maintenance contract for binding. -Each run checks phase, pane timestamps and duplicate entity states. The initial -realization uses complete bounded snapshots; partial edges, exponential -histograms and cross-run delta accumulation are unsupported. Storage identities, -revision selection, completeness/readiness evidence and scheduling stay with -deployment. +Pane construction and geometry belong to deployment. A deployment runs the +precompute DAG once per pane it constructs and binds the selected pane states to +query input slots; the query DAG merges and reads them out as computation. diff --git a/crates/asap-physical-operators/src/operators/mod.rs b/crates/asap-physical-operators/src/operators/mod.rs index 0f20fa3d..9ae0772c 100644 --- a/crates/asap-physical-operators/src/operators/mod.rs +++ b/crates/asap-physical-operators/src/operators/mod.rs @@ -21,8 +21,8 @@ mod current_series; mod filter; mod joins; mod limit; -mod panes; mod projection; +mod scope_timestamp; mod sort; mod source; mod summary; @@ -40,11 +40,6 @@ enum Kind { value: Value, dtype: DataType, }, - PaneInput { - coordinate: usize, - layout: planner_types::post_asap::PaneLayout, - offset_ms: Option, - }, ScopeTimestamp { columns: Vec>, }, @@ -260,10 +255,7 @@ impl PhysicalOperator for Operator { }; PlanProperties { boundedness, - emission: if matches!( - self.kind, - Kind::PaneInput { .. } | Kind::ScopeTimestamp { .. } - ) { + emission: if matches!(self.kind, Kind::ScopeTimestamp { .. }) { inputs .first() .map_or(Emission::Unknown, |input| input.emission) @@ -279,7 +271,6 @@ impl PhysicalOperator for Operator { match self.kind { Kind::Source(_) => "Source", Kind::Constant { .. } => "Constant", - Kind::PaneInput { .. } => "PaneInput", Kind::ScopeTimestamp { .. } => "ScopeTimestamp", Kind::Union => "Union", Kind::CurrentSeries { .. } => "CurrentSeries", @@ -303,7 +294,6 @@ impl PhysicalOperator for Operator { } } fn validate_context(&self, context: &RunContext) -> Result<(), Error> { - panes::validate_context(self, context)?; current_series::validate_context(self, context)?; self.readout_range(context).map(|_| ()) } @@ -332,9 +322,7 @@ impl PhysicalOperator for Operator { } Kind::Project(_) => projection::execute(self, inputs, context), Kind::CurrentSeries { .. } => current_series::execute(self, inputs, context), - Kind::PaneInput { .. } | Kind::ScopeTimestamp { .. } => { - panes::execute(self, inputs, context) - } + Kind::ScopeTimestamp { .. } => scope_timestamp::execute(self, inputs, context), Kind::Filter(_) => filter::execute(self, inputs, context), Kind::Limit { .. } => limit::execute(self, inputs, context), Kind::Sort { .. } => sort::execute(self, inputs, context), diff --git a/crates/asap-physical-operators/src/operators/panes.rs b/crates/asap-physical-operators/src/operators/panes.rs deleted file mode 100644 index 4748db09..00000000 --- a/crates/asap-physical-operators/src/operators/panes.rs +++ /dev/null @@ -1,228 +0,0 @@ -//! Run-scoped pane population checks and timestamp restoration after reduction. -use super::*; -use crate::runtime::Scope; -use planner_types::post_asap::{validate_pane_coverage, PaneLayout, WindowEdgeCoverage}; - -impl Operator { - pub(crate) fn pane_input( - input: Schema, - coordinate: usize, - layout: PaneLayout, - offset_ms: Option, - ) -> Result { - if plain(&input, coordinate)? != (&DataType::Timestamp, false) { - return Err(invalid("pane input requires a non-null timestamp")); - } - if layout.pane_width_ms > i64::MAX as u64 { - return Err(invalid("pane width exceeds timestamp range")); - } - validate_pane_coverage( - &layout, - layout.pane_origin_ms, - &WindowEdgeCoverage::PaneAligned, - ) - .map_err(|error| Error::Invalid(format!("invalid pane layout: {error:?}")))?; - if offset_ms.is_some_and(|offset| offset < 0) { - return Err(invalid("negative pane offset")); - } - Ok(Self { - kind: Kind::PaneInput { - coordinate, - layout, - offset_ms, - }, - inputs: vec![input.clone()], - output: input, - }) - } - - pub(crate) fn scope_timestamp(input: Schema, output: Schema) -> Result { - crate::values::validate_schema(&output)?; - let coordinate = output - .time_index - .ok_or_else(|| invalid("temporal output requires a time index"))?; - if plain(&output, coordinate)? != (&DataType::Timestamp, false) { - return Err(invalid("temporal output requires a non-null timestamp")); - } - let mut columns = Vec::new(); - let mut used = std::collections::BTreeSet::new(); - for (index, field) in output.fields.iter().enumerate() { - if index == coordinate { - columns.push(None); - continue; - } - let matches: Vec<_> = input - .fields - .iter() - .enumerate() - .filter(|(_, candidate)| { - candidate.dtype == field.dtype - && candidate.nullable == field.nullable - && (candidate.name == field.name - || !matches!(field.dtype, SummaryFamilyType::Plain(_))) - }) - .map(|(index, _)| index) - .collect(); - let [column] = matches.as_slice() else { - return Err(invalid("temporal output column missing or ambiguous")); - }; - if !used.insert(*column) { - return Err(invalid("temporal output repeats an input column")); - } - columns.push(Some(*column)); - } - if used.len() != input.fields.len() { - return Err(invalid("temporal output drops an input column")); - } - Ok(Self { - kind: Kind::ScopeTimestamp { columns }, - inputs: vec![input], - output, - }) - } -} - -pub(super) fn validate_context(operator: &Operator, context: &RunContext) -> Result<(), Error> { - let Kind::PaneInput { - layout, offset_ms, .. - } = &operator.kind - else { - return Ok(()); - }; - let end = match (&context.scope, offset_ms) { - ( - Scope::Ingestion { - window_start_ms, - window_end_ms, - .. - }, - None, - ) => { - if window_end_ms.checked_sub(*window_start_ms) != Some(layout.pane_width_ms as i64) { - return Err(invalid("maintenance run must cover exactly one pane")); - } - *window_end_ms - } - ( - Scope::Query { - evaluation_time_ms, .. - }, - Some(offset), - ) => { - let end = evaluation_time_ms - .checked_sub(*offset) - .ok_or_else(|| invalid("query pane timestamp overflows"))?; - end.checked_sub(layout.pane_width_ms as i64) - .ok_or_else(|| invalid("query pane start overflows"))?; - end - } - _ => return Err(invalid("pane operator received the wrong execution scope")), - }; - validate_pane_coverage(layout, Some(end), &WindowEdgeCoverage::PaneAligned).map_err(|error| { - Error::Invalid(format!( - "query requires aligned panes or boundary residuals: {error:?}" - )) - }) -} - -pub(super) fn execute<'a>( - operator: &'a Operator, - mut inputs: Vec>, - context: RunContext, -) -> Result, Error> { - validate_context(operator, &context)?; - let input = inputs.pop().ok_or_else(|| invalid("pane input missing"))?; - let output = operator.output.clone(); - let mut seen = std::collections::BTreeSet::new(); - let mut memory = context.reserve(0)?; - let mut key_bytes = 0; - Ok(input - .map(move |batch| { - if context.is_cancelled() { - return Err(Error::Cancelled); - } - let batch = batch?; - match &operator.kind { - Kind::PaneInput { - coordinate, - offset_ms, - .. - } => { - let groups: Vec<_> = output - .fields - .iter() - .enumerate() - .filter(|(index, field)| { - *index != *coordinate - && matches!(field.dtype, SummaryFamilyType::Plain(_)) - }) - .map(|(index, _)| index) - .collect(); - for row in batch.rows() { - let Value::Timestamp(timestamp) = row[*coordinate] else { - return Err(invalid("pane timestamp type mismatch")); - }; - match (&context.scope, offset_ms) { - ( - Scope::Ingestion { - window_start_ms, - window_end_ms, - .. - }, - None, - ) if timestamp > *window_start_ms && timestamp <= *window_end_ms => {} - ( - Scope::Query { - evaluation_time_ms, .. - }, - Some(offset), - ) if timestamp - == evaluation_time_ms - .checked_sub(*offset) - .ok_or_else(|| invalid("pane timestamp overflows"))? => - { - let key = group_key(row, &groups)?; - if seen.contains(&key) { - return Err(invalid("duplicate entity state within a pane")); - } - key_bytes += key.iter().map(Vec::len).sum::() - + key.len() * std::mem::size_of::>() - + 64; - memory.resize(key_bytes)?; - seen.insert(key); - } - _ => { - return Err(invalid("input population differs from required pane")) - } - } - } - Ok(batch.value().clone()) - } - Kind::ScopeTimestamp { columns } => { - let timestamp = match context.scope { - Scope::Ingestion { window_end_ms, .. } => window_end_ms, - Scope::Query { - evaluation_time_ms, .. - } => evaluation_time_ms, - }; - let rows = batch - .rows() - .iter() - .map(|row| { - columns - .iter() - .map(|column| { - column.map_or(Value::Timestamp(timestamp), |column| { - row[column].clone() - }) - }) - .collect() - }) - .collect(); - Batch::try_new(output.clone(), rows) - } - _ => unreachable!(), - } - }) - .boxed_local()) -} diff --git a/crates/asap-physical-operators/src/operators/scope_timestamp.rs b/crates/asap-physical-operators/src/operators/scope_timestamp.rs new file mode 100644 index 00000000..3561d6a6 --- /dev/null +++ b/crates/asap-physical-operators/src/operators/scope_timestamp.rs @@ -0,0 +1,91 @@ +//! Run-scoped timestamp restoration after reduction. +use super::*; +use crate::runtime::Scope; + +impl Operator { + pub(crate) fn scope_timestamp(input: Schema, output: Schema) -> Result { + crate::values::validate_schema(&output)?; + let coordinate = output + .time_index + .ok_or_else(|| invalid("temporal output requires a time index"))?; + if plain(&output, coordinate)? != (&DataType::Timestamp, false) { + return Err(invalid("temporal output requires a non-null timestamp")); + } + let mut columns = Vec::new(); + let mut used = std::collections::BTreeSet::new(); + for (index, field) in output.fields.iter().enumerate() { + if index == coordinate { + columns.push(None); + continue; + } + let matches: Vec<_> = input + .fields + .iter() + .enumerate() + .filter(|(_, candidate)| { + candidate.dtype == field.dtype + && candidate.nullable == field.nullable + && (candidate.name == field.name + || !matches!(field.dtype, SummaryFamilyType::Plain(_))) + }) + .map(|(index, _)| index) + .collect(); + let [column] = matches.as_slice() else { + return Err(invalid("temporal output column missing or ambiguous")); + }; + if !used.insert(*column) { + return Err(invalid("temporal output repeats an input column")); + } + columns.push(Some(*column)); + } + if used.len() != input.fields.len() { + return Err(invalid("temporal output drops an input column")); + } + Ok(Self { + kind: Kind::ScopeTimestamp { columns }, + inputs: vec![input], + output, + }) + } +} + +pub(super) fn execute<'a>( + operator: &'a Operator, + mut inputs: Vec>, + context: RunContext, +) -> Result, Error> { + let Kind::ScopeTimestamp { columns } = &operator.kind else { + return Err(invalid("scope timestamp operator required")); + }; + let input = inputs + .pop() + .ok_or_else(|| invalid("scope timestamp input missing"))?; + let output = operator.output.clone(); + let timestamp = match context.scope { + Scope::Ingestion { window_end_ms, .. } => window_end_ms, + Scope::Query { + evaluation_time_ms, .. + } => evaluation_time_ms, + }; + Ok(input + .map(move |batch| { + if context.is_cancelled() { + return Err(Error::Cancelled); + } + let batch = batch?; + let rows = batch + .rows() + .iter() + .map(|row| { + columns + .iter() + .map(|column| { + column.map_or(Value::Timestamp(timestamp), |column| row[column].clone()) + }) + .collect() + }) + .collect(); + Batch::try_new(output.clone(), rows) + }) + .boxed_local()) +} diff --git a/crates/asap-physical-operators/src/operators/unchecked.rs b/crates/asap-physical-operators/src/operators/unchecked.rs index 2ecbaabe..f927a4e4 100644 --- a/crates/asap-physical-operators/src/operators/unchecked.rs +++ b/crates/asap-physical-operators/src/operators/unchecked.rs @@ -30,11 +30,6 @@ impl TryFrom for Operator { let op = match kind { Kind::Source(_) => return Err(invalid("physical plans cannot serialize live sources")), Kind::Constant { value, dtype } => Operator::scalar(value, dtype)?, - Kind::PaneInput { - coordinate, - layout, - offset_ms, - } => Operator::pane_input(input(0)?, coordinate, layout, offset_ms)?, Kind::ScopeTimestamp { .. } => Operator::scope_timestamp(input(0)?, output.clone())?, Kind::Union => Operator::union(input(0)?, inputs.len())?, Kind::CurrentSeries { diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index 142a2c25..112e1715 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -40,12 +40,6 @@ pub use candidates::{ CandidateSelection, PhysicalCandidate, }; -mod temporal_panes; -pub use temporal_panes::{ - compile_temporal_pane_candidate, TemporalEntityIdentity, TemporalPaneCandidate, - TemporalPaneMaintenance, -}; - mod compiled; pub use compiled::{CompiledPhysicalDag, InputContract}; diff --git a/crates/asap-physical-operators/src/physical_planner/temporal_panes.rs b/crates/asap-physical-operators/src/physical_planner/temporal_panes.rs deleted file mode 100644 index d535c777..00000000 --- a/crates/asap-physical-operators/src/physical_planner/temporal_panes.rs +++ /dev/null @@ -1,331 +0,0 @@ -//! Lower a selected temporal maintenance contract; deployment supplies readers. -use super::*; -use planner_types::post_asap::{ - EvaluationSchedule, OutputRepresentation, PaneLayout, SketchAlgorithm, - SummaryMaintenanceLifecycle, SummaryMaintenanceLifecycleGuarantee, SummaryMaintenanceMode, - SummaryWindowFramework, -}; - -/// Resolved source identity, supplied with physical capability evidence. -/// A schemaless PromQL projection cannot establish the complete label set. -#[derive(Clone, Debug)] -pub enum TemporalEntityIdentity { - /// The input resolver guarantees that the slot contains one entity. - SingleEntity, - /// All entity keys are represented by these columns; there are no hidden - /// labels distinguishing two rows with the same key. - Columns(Vec), -} - -/// Planner-selected lifecycle/window requirements for one temporal producer. -/// Pane geometry is semantic input, not a storage identity or scheduling policy. -#[derive(Clone, Debug)] -pub struct TemporalPaneMaintenance { - pub summary_node: NodeId, - pub lifecycle: SummaryMaintenanceLifecycleGuarantee, - pub framework: SummaryWindowFramework, - pub layout: PaneLayout, - pub entity_identity: TemporalEntityIdentity, -} - -/// Generated precompute and query computation. `pane_inputs` is ordered from -/// the oldest complete pane to the newest; each run checks actual timestamps. -#[derive(Clone)] -pub struct TemporalPaneCandidate { - pub physical: PhysicalCandidate, - pub maintenance: TemporalPaneMaintenance, - pub pane_inputs: Vec, - pub merged_state: NodeId, - pub window_width_ms: u64, -} - -/// Compile bounded pane construction and a shared pane merge for temporal KLL -/// quantile roots. The selected contract remains authoritative; unsupported -/// lifecycle/framework/operator shapes fail rather than being substituted. -/// This initial realization consumes complete pane populations and emits full -/// state snapshots. Cross-run delta accumulation belongs to other candidates. -pub fn compile_temporal_pane_candidate( - dag: &PostAsapDag, - inputs: BTreeMap, - roots: &[NodeId], - maintenance: &TemporalPaneMaintenance, -) -> Result { - dag.validate().map_err(|error| invalid(error.to_string()))?; - if maintenance.lifecycle.summary_maintenance_lifecycle - != SummaryMaintenanceLifecycle::ContinuouslyMaintained - || maintenance.lifecycle.summary_maintenance_mode != SummaryMaintenanceMode::Incremental - || maintenance.lifecycle.evaluation_schedule != EvaluationSchedule::PerUpdate - || maintenance.lifecycle.output_representation != OutputRepresentation::SummaryState - { - return Err(invalid( - "pane candidate requires continuous incremental summary maintenance", - )); - } - let build = dag - .nodes - .iter() - .find(|node| u64::from(node.id.0) == maintenance.summary_node) - .ok_or_else(|| invalid("unknown maintained producer"))?; - let Payload::SummaryAgg { - family, - input: update, - reduction: PlannerReduction::PerEntity, - grouping, - } = &build.payload - else { - return Err(invalid( - "pane candidate requires a temporal per-entity summary", - )); - }; - if !matches!(family, SummaryFamilyType::Sketch(kind, _) if kind.algorithm() == &SketchAlgorithm::Kll) - || update.item.is_some() - { - return Err(invalid("pane candidate supports unkeyed temporal KLL only")); - } - crate::capability::validate_summary_kernel(family, update, grouping).map_err(Error::Invalid)?; - let dependencies: Vec<_> = dag - .edges - .iter() - .filter(|edge| edge.consumer == build.id) - .map(|edge| edge.producer) - .collect(); - let [raw_id] = dependencies.as_slice() else { - return Err(invalid("temporal producer requires one raw input")); - }; - let raw = dag - .nodes - .iter() - .find(|node| node.id == *raw_id) - .ok_or_else(|| invalid("missing raw input"))?; - let Payload::Fallback { - expression: QueryExpr::TimeRange { range, child }, - } = &raw.payload - else { - return Err(invalid( - "temporal producer requires an explicit logical time range", - )); - }; - let QueryExpr::Scan { predicates, .. } = child.as_ref() else { - return Err(invalid("temporal pane source requires a raw scan")); - }; - let window_width_ms: u64 = range - .as_millis() - .try_into() - .map_err(|_| invalid("temporal window overflows"))?; - if window_width_ms == 0 - || window_width_ms > i64::MAX as u64 - || range.subsec_nanos() % 1_000_000 != 0 - { - return Err(invalid( - "temporal window requires positive integral milliseconds", - )); - } - let width = maintenance.layout.pane_width_ms; - if width == 0 || width > window_width_ms || !window_width_ms.is_multiple_of(width) { - return Err(invalid("temporal window must contain whole panes")); - } - match maintenance.framework { - SummaryWindowFramework::Sliding => {} - SummaryWindowFramework::Tumbling if width == window_width_ms => {} - _ => return Err(invalid("unsupported temporal window realization")), - } - let count = window_width_ms / width; - if count > 4096 { - return Err(invalid("temporal pane candidate exceeds input budget")); - } - let raw_id = u64::from(raw_id.0); - if inputs.len() != 1 { - return Err(invalid( - "pane candidate requires exactly its raw input contract", - )); - } - let contract = inputs - .get(&raw_id) - .ok_or_else(|| invalid("missing raw input contract"))?; - let raw_schema = Arc::new(raw.output_schema.clone()); - if contract.schema != raw_schema || contract.properties.boundedness != Boundedness::Bounded { - return Err(invalid("pane source requires its declared bounded schema")); - } - let coordinate = raw_schema - .time_index - .ok_or_else(|| invalid("temporal source requires a time index"))?; - let SummaryInputExpr::Column(value) = &update.weight else { - return Err(invalid("pane builder requires a value column")); - }; - let value = named_column(&raw_schema, value)?; - let groups: Vec<_> = (0..raw_schema.fields.len()) - .filter(|&index| index != coordinate && index != value) - .collect(); - match &maintenance.entity_identity { - TemporalEntityIdentity::SingleEntity if groups.is_empty() => {} - TemporalEntityIdentity::Columns(columns) - if !columns.is_empty() - && columns.len() == columns.iter().collect::>().len() - && columns.iter().copied().collect::>() - == groups.iter().copied().collect() => {} - _ => { - return Err(invalid( - "pane input requires its complete resolved entity identity", - )) - } - } - let mut next = dag - .nodes - .iter() - .map(|node| u64::from(node.id.0)) - .max() - .unwrap_or(0) - + 1; - let mut allocate = || { - let id = next; - next += 1; - id - }; - let mut operators = BTreeMap::new(); - let guard = allocate(); - operators.insert( - guard, - ( - vec![raw_id], - Operator::pane_input( - raw_schema.clone(), - coordinate, - maintenance.layout.clone(), - None, - )?, - ), - ); - let mut previous = guard; - for predicate in predicates { - let id = allocate(); - operators.insert( - id, - ( - vec![previous], - Operator::filter(raw_schema.clone(), expression(&predicate.0, &raw_schema)?)?, - ), - ); - previous = id; - } - let native = - Operator::summary_build(raw_schema, family.clone(), value, Some(coordinate), groups)?; - let compact_state = native.schema(); - let native_id = allocate(); - operators.insert(native_id, (vec![previous], native)); - let state_schema = Arc::new(build.output_schema.clone()); - let pane_output = allocate(); - operators.insert( - pane_output, - ( - vec![native_id], - Operator::scope_timestamp(compact_state, state_schema.clone())?, - ), - ); - let precompute = CompiledPhysicalDag::from_operators(inputs, operators, vec![pane_output])?; - let state_coordinate = state_schema - .time_index - .ok_or_else(|| invalid("pane state requires a time index"))?; - let state_column = summary_column(&state_schema)?; - let mut query_inputs = BTreeMap::new(); - let mut operators = BTreeMap::new(); - let mut pane_inputs = Vec::new(); - let mut guarded_inputs = Vec::new(); - for pane in 0..count { - let input = allocate(); - let guard = allocate(); - query_inputs.insert(input, InputContract::bounded(state_schema.clone())); - let offset = ((count - 1 - pane) * width) as i64; - operators.insert( - guard, - ( - vec![input], - Operator::pane_input( - state_schema.clone(), - state_coordinate, - maintenance.layout.clone(), - Some(offset), - )?, - ), - ); - pane_inputs.push(input); - guarded_inputs.push(guard); - } - let union = allocate(); - operators.insert( - union, - ( - guarded_inputs, - Operator::union(state_schema.clone(), count as usize)?, - ), - ); - let merge = Operator::summary_merge( - state_schema.clone(), - state_column, - (0..state_schema.fields.len()) - .filter(|&index| index != state_coordinate && index != state_column) - .collect(), - )?; - let merged_schema = merge.schema(); - let merged_state = allocate(); - operators.insert(merged_state, (vec![union], merge)); - if roots.is_empty() || roots.iter().copied().collect::>().len() != roots.len() { - return Err(invalid("temporal query requires distinct output roots")); - } - for &root in roots { - let node = dag - .nodes - .iter() - .find(|node| u64::from(node.id.0) == root) - .ok_or_else(|| invalid("unknown temporal output root"))?; - let Payload::SummaryEstimate { - query: SketchQuery::Quantile { q }, - } = node.payload - else { - return Err(invalid("temporal pane root must be a KLL quantile")); - }; - if !q.is_finite() || !(0. ..=1.).contains(&q) { - return Err(invalid("invalid temporal quantile")); - } - let dependencies: Vec<_> = dag - .edges - .iter() - .filter(|edge| edge.consumer == node.id) - .map(|edge| u64::from(edge.producer.0)) - .collect(); - if dependencies != [maintenance.summary_node] { - return Err(invalid( - "temporal readout must consume the maintained producer", - )); - } - let readout = Operator::readout( - merged_schema.clone(), - summary_column(&merged_schema)?, - ReadoutQuery::Sketch(SketchQuery::Quantile { q }), - )?; - let readout_schema = readout.schema(); - let readout_id = allocate(); - operators.insert(readout_id, (vec![merged_state], readout)); - operators.insert( - root, - ( - vec![readout_id], - Operator::scope_timestamp(readout_schema, Arc::new(node.output_schema.clone()))?, - ), - ); - } - let query = CompiledPhysicalDag::from_operators(query_inputs, operators, roots.to_vec())?; - let mut output = precompute.output_contract(pane_output)?; - // Persisted readers have independent timing from the blocking builder. - output.properties.emission = Emission::Unknown; - Ok(TemporalPaneCandidate { - physical: PhysicalCandidate { - precompute: Some(precompute), - query, - materialized_outputs: BTreeMap::from([(pane_output, output)]), - }, - maintenance: maintenance.clone(), - pane_inputs, - merged_state, - window_width_ms, - }) -} diff --git a/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs b/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs index 41e5f1b3..d5406f23 100644 --- a/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs +++ b/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs @@ -426,446 +426,3 @@ fn continuous_lifecycle_compiles_and_executes_spatial_kll() { ); } } - -struct SlidingPaneModel; -impl CostModel for SlidingPaneModel { - fn raw_query_recompute_total_cost( - &self, - target: &asap_types::pre_asap::QueryExpr, - reads: f64, - ) -> Option { - let _ = (target, reads); - Some(Cost(100_000.0)) - } - fn rank_candidates( - &self, - intent: &AggIntent, - candidates: &[asap_types::post_asap::SketchAlgorithm], - ) -> Vec { - FullyCostedRuntime.rank_candidates(intent, candidates) - } - fn summary_maintenance_lifecycle_cost_inputs( - &self, - summary: &SummaryNode, - ) -> SummaryMaintenanceLifecycleCostInputs { - let mut costs = FullyCostedRuntime.summary_maintenance_lifecycle_cost_inputs(summary); - // Controlled workload evidence makes repeated raw construction more - // expensive than retaining and updating the same temporal population. - costs.build_cost = Some(Cost(1000.)); - costs - } - fn summary_maintenance_capabilities( - &self, - summary: &SummaryNode, - ) -> SummaryMaintenanceCapabilities { - FullyCostedRuntime.summary_maintenance_capabilities(summary) - } - fn complete_summary_candidate_estimate( - &self, - _root: &SummaryNode, - _target: Option<&asap_types::pre_asap::QueryExpr>, - deployments: &[asap_aware_mapping::cost_model::CostedSummaryDeployment<'_>], - _horizon: Option, - _reads: Option, - _accuracy: &[AccuracyTarget], - ) -> Option { - Some(asap_aware_mapping::CompleteSummaryCandidateEstimate { - cost: Cost( - deployments - .iter() - .map(|deployment| deployment.selected_cost.0) - .sum(), - ), - physical_plan_id: Some("bounded-sliding-pane-evidence".into()), - window_frameworks: deployments - .iter() - .map(|deployment| { - (deployment.guarantee.summary_maintenance_lifecycle - == SummaryMaintenanceLifecycle::ContinuouslyMaintained) - .then_some(asap_types::post_asap::SummaryWindowFramework::Sliding) - }) - .collect(), - window_accuracy_guarantee: None, - }) - } -} - -/// Workload and optimizer-selected lifecycle generate both physical DAGs. -/// No computational operators or graph edges are constructed by this fixture. -#[test] -fn selected_temporal_lifecycle_compiles_panes_and_executes() { - use asap_physical_operators::{ - operators::Operator, - physical_planner::{ - compile_temporal_pane_candidate, InputContract, Source, TemporalEntityIdentity, - TemporalPaneMaintenance, - }, - runtime::{Limits, RunContext, Scope}, - summary_kernels::datasketches_kll::DatasketchesKLLAccumulator, - values::{Batch, Value}, - }; - use asap_types::{ - post_asap::{ - compile_post_asap_dag, plan_pane_phase, PostAsapOperatorPayload, SummaryFamilyType, - SummaryWindowFramework, - }, - pre_asap::DataType, - workload::TimestampMs, - }; - use futures::{executor::block_on, StreamExt}; - use std::{collections::BTreeMap, sync::Arc}; - for quantile in [0.5, 0.99] { - let mut workload = dashboard_workload(); - let query = Query(format!( - "quantile_over_time({quantile}, latency{{job=\"api\"}}[5m])" - )); - workload.query_workload.query_batch.as_mut().unwrap()[0].query = query.clone(); - workload.query_workload.repeating_queries.as_mut().unwrap()[0].query = query; - - workload - .data_workload - .as_mut() - .unwrap() - .data_ingestion_interval - .value = Some(DurationMs(60_000)); - workload.query_workload.repeating_queries.as_mut().unwrap()[0].demand = - RepeatedDemand::FixedIntervalAt { - interval: RepetitionInterval(60_000), - evaluation_phase: TimestampMs(300_000), - }; - let plan = selected_plan_with_horizon(&workload, &SlidingPaneModel, Horizon(1000.)); - assert!(!plan.selected_raw_recompute); - assert_eq!(plan.deployments.len(), 1); - let deployment = &plan.deployments[0]; - assert_eq!( - deployment.selected_window_framework, - Some(SummaryWindowFramework::Sliding) - ); - let dag = compile_post_asap_dag(&plan.root).unwrap(); - let build = dag - .nodes - .iter() - .find(|node| matches!(node.payload, PostAsapOperatorPayload::SummaryAgg { .. })) - .unwrap(); - assert_eq!(build.id, deployment.post_asap_node_id); - let raw = dag - .nodes - .iter() - .find(|node| matches!(node.payload, PostAsapOperatorPayload::Fallback { .. })) - .unwrap(); - let schema = Arc::new(raw.output_schema.clone()); - let width = workload - .data_workload - .as_ref() - .unwrap() - .data_ingestion_interval - .value - .unwrap() - .0; - let layout = plan_pane_phase( - &workload.query_workload.repeating_queries.as_ref().unwrap()[0].demand, - width, - ) - .unwrap(); - let maintenance = TemporalPaneMaintenance { - summary_node: u64::from(build.id.0), - lifecycle: deployment - .summary_maintenance_lifecycle_guarantee - .clone() - .unwrap(), - framework: deployment.selected_window_framework.clone().unwrap(), - layout, - // The memory source has exactly the declared label columns; a - // schemaless deployment must resolve all entity keys first. - entity_identity: TemporalEntityIdentity::Columns( - schema - .fields - .iter() - .enumerate() - .filter(|(_, field)| field.name == "job") - .map(|(index, _)| index) - .collect(), - ), - }; - let candidate = compile_temporal_pane_candidate( - &dag, - BTreeMap::from([(u64::from(raw.id.0), InputContract::bounded(schema.clone()))]), - &[u64::from(dag.root.0)], - &maintenance, - ) - .unwrap(); - assert_eq!(candidate.window_width_ms, 300_000); - assert_eq!(candidate.pane_inputs.len(), 5); - assert_ne!( - candidate.physical.precompute.as_ref().unwrap().roots()[0], - maintenance.summary_node, - "a one-minute pane is not the logical five-minute summary output" - ); - let source_opens = Arc::new(std::sync::atomic::AtomicUsize::new(0)); - let mut stored = Vec::new(); - for pane in 0..6 { - let rows = (0..20) - .flat_map(|sample| { - ["api", "batch"] - .into_iter() - .map(move |entity| (sample, entity)) - }) - .map(|(sample, entity)| { - schema - .fields - .iter() - .map(|field| match field.dtype { - SummaryFamilyType::Plain(DataType::Timestamp) => { - Value::Timestamp(pane * 60_000 + (sample + 1) * 3000) - } - SummaryFamilyType::Plain(DataType::Float64) => Value::Float64( - (pane * 20 + sample) as f64 - + if entity == "batch" { 100_000. } else { 0. }, - ), - SummaryFamilyType::Plain(DataType::Utf8) => Value::Utf8(entity.into()), - _ => panic!("unexpected raw field {field:?}"), - }) - .collect() - }) - .collect(); - let result = physical_common::execute( - candidate.physical.precompute.as_ref().unwrap(), - BTreeMap::from([( - u64::from(raw.id.0), - Batch::try_new(schema.clone(), rows).unwrap(), - )]), - Scope::Ingestion { - window_start_ms: pane * 60_000, - window_end_ms: (pane + 1) * 60_000, - revision: 1, - }, - ); - let batch = &result[0][0]; - stored.push(batch.clone()); - } - for offset in [0, 1] { - let inputs: BTreeMap<_, _> = candidate - .pane_inputs - .iter() - .enumerate() - .map(|(index, &id)| (id, stored[offset + index].clone())) - .collect(); - let scope = Scope::Query { - evaluation_time_ms: (5 + offset as i64) * 60_000, - revision: 1, - }; - let result = - physical_common::execute(&candidate.physical.query, inputs.clone(), scope.clone()); - let rows = result[0][0].rows(); - assert_eq!(rows.len(), 1); - assert!(rows[0] - .iter() - .any(|value| matches!(value, Value::Utf8(label) if label.as_ref() == "api"))); - let Value::Float64(value) = rows[0][1] else { - panic!("missing p99") - }; - assert!( - (value - ((if quantile == 0.5 { 50 } else { 99 }) + offset * 20) as f64).abs() - <= 1. - ); - assert!( - matches!(rows[0][0], Value::Timestamp(timestamp) if timestamp == (5 + offset as i64) * 60_000) - ); - let sources = |inputs: BTreeMap| -> BTreeMap> { - inputs - .into_iter() - .map(|(id, batch)| { - ( - id, - Box::new(TemporalCountingSource { - operator: Operator::source(batch.schema().clone(), vec![batch]) - .unwrap(), - opens: source_opens.clone(), - }) as Source<'static>, - ) - }) - .collect() - }; - let bound = candidate - .physical - .query - .instantiate(sources(inputs.clone())) - .unwrap(); - let population = block_on( - bound - .execute( - &[candidate.merged_state], - RunContext::new(scope.clone(), Limits::default()).unwrap(), - ) - .unwrap() - .remove(0) - .collect::>(), - ); - let Value::Summary { state, .. } = population[0].as_ref().unwrap().rows()[0] - .iter() - .find(|value| matches!(value, Value::Summary { .. })) - .unwrap() - else { - panic!("missing merged state") - }; - assert_eq!( - state - .as_any() - .downcast_ref::() - .unwrap() - .inner - .count(), - 100 - ); - let mut missing = inputs.clone(); - missing.remove(&candidate.pane_inputs[0]); - assert!(candidate - .physical - .query - .instantiate(sources(missing)) - .is_err()); - let mut duplicate = inputs.clone(); - duplicate.insert(candidate.pane_inputs[1], stored[offset].clone()); - let bad = candidate - .physical - .query - .instantiate(sources(duplicate)) - .unwrap(); - let errors = block_on( - bad.execute( - candidate.physical.query.roots(), - RunContext::new(scope.clone(), Limits::default()).unwrap(), - ) - .unwrap() - .remove(0) - .collect::>(), - ); - assert!( - errors.iter().any(Result::is_err), - "duplicate pane must not be merged twice" - ); - let mut duplicate_entity = inputs.clone(); - let pane = &stored[offset]; - duplicate_entity.insert( - candidate.pane_inputs[0], - Batch::try_new( - pane.schema().clone(), - vec![pane.rows()[0].clone(), pane.rows()[0].clone()], - ) - .unwrap(), - ); - let bad = candidate - .physical - .query - .instantiate(sources(duplicate_entity)) - .unwrap(); - let errors = block_on( - bad.execute( - candidate.physical.query.roots(), - RunContext::new(scope.clone(), Limits::default()).unwrap(), - ) - .unwrap() - .remove(0) - .collect::>(), - ); - assert!( - errors.iter().any(Result::is_err), - "duplicate snapshots within a pane must fail" - ); - let bound = candidate - .physical - .query - .instantiate(sources(inputs)) - .unwrap(); - source_opens.store(0, std::sync::atomic::Ordering::SeqCst); - assert!(bound - .execute( - candidate.physical.query.roots(), - RunContext::new( - Scope::Query { - evaluation_time_ms: 330_000, - revision: 1 - }, - Limits::default() - ) - .unwrap() - ) - .is_err()); - assert_eq!( - source_opens.load(std::sync::atomic::Ordering::SeqCst), - 0, - "invalid phase must fail before opening readers" - ); - } - // A selected framework cannot be silently replaced by physical planning. - let mut wrong_framework = maintenance.clone(); - wrong_framework.framework = SummaryWindowFramework::ExponentialHistogram; - let mut wrong_identity = maintenance.clone(); - wrong_identity.entity_identity = TemporalEntityIdentity::SingleEntity; - let mut unknown_phase = maintenance.clone(); - unknown_phase.layout.pane_origin_ms = None; - let mut partial_panes = maintenance.clone(); - partial_panes.layout.pane_width_ms = 90_000; - let mut wrong_lifecycle = maintenance.clone(); - wrong_lifecycle.lifecycle.summary_maintenance_lifecycle = - SummaryMaintenanceLifecycle::Ephemeral; - for unsupported in [ - wrong_framework, - wrong_identity, - unknown_phase, - partial_panes, - wrong_lifecycle, - ] { - assert!(compile_temporal_pane_candidate( - &dag, - BTreeMap::from([(u64::from(raw.id.0), InputContract::bounded(schema.clone()))]), - &[u64::from(dag.root.0)], - &unsupported - ) - .is_err()); - } - } -} - -struct TemporalCountingSource { - operator: asap_physical_operators::operators::Operator, - opens: std::sync::Arc, -} -impl - asap_physical_operators::plan::PhysicalOperator< - asap_physical_operators::values::Batch, - asap_physical_operators::values::Schema, - > for TemporalCountingSource -{ - fn name(&self) -> &str { - "TemporalCountingSource" - } - fn properties( - &self, - inputs: &[asap_physical_operators::plan::PlanProperties], - ) -> asap_physical_operators::plan::PlanProperties { - self.operator.properties(inputs) - } - fn input_schemas(&self) -> Vec { - self.operator.input_schemas() - } - fn output_schema(&self) -> asap_physical_operators::values::Schema { - self.operator.output_schema() - } - fn output_bytes(&self, batch: &asap_physical_operators::values::Batch) -> usize { - batch.bytes() - } - fn start<'a>( - &'a self, - inputs: Vec< - asap_physical_operators::runtime::Input<'a, asap_physical_operators::values::Batch>, - >, - context: asap_physical_operators::runtime::RunContext, - ) -> Result< - asap_physical_operators::runtime::OutputStream<'a, asap_physical_operators::values::Batch>, - asap_physical_operators::Error, - > { - self.opens.fetch_add(1, std::sync::atomic::Ordering::SeqCst); - self.operator.start(inputs, context) - } -} diff --git a/docs/design_docs/physical-planning-and-deployment.md b/docs/design_docs/physical-planning-and-deployment.md index 72313a9d..d9c91009 100644 --- a/docs/design_docs/physical-planning-and-deployment.md +++ b/docs/design_docs/physical-planning-and-deployment.md @@ -295,7 +295,7 @@ NativeKllBuild(k=200) KllStateOutput(k=200) ``` -This DAG implements construction of each maintained one-minute pane. Its input +This DAG computes the state of one maintained one-minute pane. Its input contract requires all input samples matching the source, filters and group within that pane; the deployment supplies that bounded input from its source integration. @@ -480,33 +480,18 @@ snapshots as separate inputs. ## 6. Executable acceptance coverage -The tests cover optimizer-selected lifecycle execution and automatic temporal -pane compilation, alongside independent operator/runtime fixtures: +The tests cover optimizer-selected lifecycle execution alongside independent +operator/runtime fixtures: | Test | Contract exercised | | --- | --- | -| `summary_maintenance_lifecycle_e2e::selected_temporal_lifecycle_compiles_panes_and_executes` | PromQL p50/p99 workloads → selected continuous lifecycle and Sliding framework → automatically generated precompute/query DAGs → real codec round-trip → adjacent aligned windows; checks filters, entity identity, sample counts, missing/duplicate panes and phase rejection before opening readers | | `summary_maintenance_lifecycle_e2e::continuous_lifecycle_compiles_and_executes_spatial_kll` | PromQL workload → selected continuous lifecycle → logical DAG → compiled precompute/query candidate → results in independent revisions; an unbounded candidate fails before pricing, and a bounded request candidate summarizes the same input samples | | `kll_pane_execution::five_panes_roundtrip_and_shared_merge_runs_once` | Explicit one-minute precompute DAGs → real MessagePack state bytes → five required query inputs → shared native merge → p50/p99; counts every sample once, checks adjacent aligned windows and instruments one merge start per run | | `kll_pane_execution::restored_panes_reject_corruption_parameters_schema_and_missing_binding` | Corrupt bytes, parameter relabelling, incompatible schemas and absent bindings fail explicitly | | `precompute_candidates::grouped_rate_can_be_materialized_before_or_after_grouped_sum` | Cost changes select different legal precompute frontiers; both selected candidates execute with the same reset-sensitive result; uncompilable candidates are not priced | | `sql_to_physical::sql_filter_grouped_sum_executes_and_rebinds` | SQL text → candidate search → physical compilation → shared Scan predicates and grouped summary execution; NULL samples are ignored and fresh bindings produce new results | -`physical_planner::compile_temporal_pane_candidate` consumes the logical DAG, -selected lifecycle/framework and a generic pane/entity input contract. It -generates pane construction, scan predicates, ordered state slots, a shared -merge, quantile readouts and run-scoped timestamps. A physical pane output has -its own identity: one minute of state cannot masquerade as the logical -five-minute summary. The returned candidate retains the maintenance contract. - -This initial realization supports bounded, complete KLL panes with known phase -and resolved entity identity, for Sliding windows or a single Tumbling window. -Source capability evidence must declare all entity keys or isolate one entity; -usage-derived PromQL columns alone cannot establish that identity. Partial edge -panes, exponential histograms and cross-run delta accumulation require further -physical candidates and are rejected by this entry point. - -Physical execution checks pane timestamps and duplicate entity states. Concrete -stored identity, revisions, readiness and complete coverage of required input samples remain -deployment responsibilities. Real storage and HTTP execution belong to +Pane construction, pane timestamp checks, stored identity, revisions, readiness +and complete coverage of required input samples are deployment +responsibilities. Real storage and HTTP execution belong to deployment-repository E2E tests. From fb92998a97bd8910d104023c8ea6aa2e16890eab Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 02:23:03 +0000 Subject: [PATCH 2/2] docs: state what the physical layer does not own Co-Authored-By: Claude Opus 5.5 --- docs/design_docs/physical-planning-and-deployment.md | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/docs/design_docs/physical-planning-and-deployment.md b/docs/design_docs/physical-planning-and-deployment.md index d9c91009..fe71c8c3 100644 --- a/docs/design_docs/physical-planning-and-deployment.md +++ b/docs/design_docs/physical-planning-and-deployment.md @@ -474,6 +474,13 @@ shared physical operator implementation library, `asap-physical-operators`, and its DAG runtime. The merge executes once per run for both consumers. Execution does not introduce additional planning decisions. +The physical layer does not own raw ingestion, pane construction or geometry, +storage formats, or decoding persisted bytes into typed state. It compiles +computation over typed input contracts: summary build, merge (for example KLL +merge), sketch estimates and exact finalization. The deployment constructs panes, +reads and decodes stored state, and binds the typed values to input slots. +Compiled physical plans are Planner outputs and keep their own serialized form. + Each maintained pane contributes its input samples once. A replacement snapshot replaces that pane's state; query merging must not count both the old and new snapshots as separate inputs.