diff --git a/Cargo.lock b/Cargo.lock index 532bb94df..776d32ec2 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -373,7 +373,7 @@ dependencies = [ [[package]] name = "asap-aware-mapping" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=e78572c57e6c3eea2a9017f968fbc65a9196d56e#e78572c57e6c3eea2a9017f968fbc65a9196d56e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21#1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21" dependencies = [ "asap-types", "serde", @@ -384,7 +384,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=e78572c57e6c3eea2a9017f968fbc65a9196d56e#e78572c57e6c3eea2a9017f968fbc65a9196d56e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21#1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21" dependencies = [ "asap-types", "promql-parser", @@ -393,7 +393,7 @@ dependencies = [ [[package]] name = "asap-frontend-sql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=e78572c57e6c3eea2a9017f968fbc65a9196d56e#e78572c57e6c3eea2a9017f968fbc65a9196d56e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21#1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -415,12 +415,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=e78572c57e6c3eea2a9017f968fbc65a9196d56e#e78572c57e6c3eea2a9017f968fbc65a9196d56e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21#1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=e78572c57e6c3eea2a9017f968fbc65a9196d56e#e78572c57e6c3eea2a9017f968fbc65a9196d56e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21#1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21" dependencies = [ "serde", "serde_json", diff --git a/control_plane/Cargo.toml b/control_plane/Cargo.toml index 6f928922b..eb0032b1a 100644 --- a/control_plane/Cargo.toml +++ b/control_plane/Cargo.toml @@ -76,8 +76,8 @@ asap_types.workspace = true # scaffolding, unaware that `data_plane`'s `summary_executor.rs` in *this* # repo is a real one. Vendored locally instead of chased upstream -- see # `data_plane/src/query_engines/asap_query_engine/summary_exec.rs`. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "e78572c57e6c3eea2a9017f968fbc65a9196d56e" } -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "e78572c57e6c3eea2a9017f968fbc65a9196d56e" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21" } # L1 adoption (design-target-architecture.md Part B): the PromQL front # end itself, replacing control_plane's own query_parser/promql.rs. @@ -85,8 +85,8 @@ asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = # `planner-types`/`asap-aware-mapping` above -- these three MUST move # together (two revs of the same upstream repo's types in one workspace # resolve to distinct Rust types that won't unify). -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "e78572c57e6c3eea2a9017f968fbc65a9196d56e" } -asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "e78572c57e6c3eea2a9017f968fbc65a9196d56e" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21" } +asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21" } [dev-dependencies] tokio = { version = "1", features = ["full", "test-util"] } diff --git a/control_plane/examples/calibration_candidates.rs b/control_plane/examples/calibration_candidates.rs index 3b119bb70..65fcb1369 100644 --- a/control_plane/examples/calibration_candidates.rs +++ b/control_plane/examples/calibration_candidates.rs @@ -27,10 +27,15 @@ fn planner_forest(queries: &[control_plane::physical::compiler::PlanningQuery]) vec![], json!({"query_expr_debug":format!("{expr:#?}")}), ), - SummaryExpr::BinaryOp { lhs, rhs, operator } => ( + SummaryExpr::BinaryOp { + lhs, + rhs, + operator, + timing, + } => ( "BinaryOp", vec![lhs, rhs], - json!({"operator_debug":format!("{operator:?}")}), + json!({"operator_debug":format!("{operator:?}"),"timing_debug":format!("{timing:?}")}), ), SummaryExpr::CandidateTopK { candidates, diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index b3788499d..9fb0a03e3 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -4602,6 +4602,7 @@ mod tests { let selected = request.queries[0].post_asap.clone(); request.queries[0].post_asap = Rc::new(SummaryNode { expr: SummaryExpr::BinaryOp { + timing: planner_types::post_asap::ExecutionTiming::ReadTime, lhs: selected.clone(), rhs: selected.clone(), operator: planner_types::post_asap::BinaryOperator { diff --git a/control_plane/src/query_plan.rs b/control_plane/src/query_plan.rs index 34fcdbae2..475e319a0 100644 --- a/control_plane/src/query_plan.rs +++ b/control_plane/src/query_plan.rs @@ -606,7 +606,12 @@ where completeness: completeness.clone(), } } - SummaryExpr::BinaryOp { lhs, rhs, operator } if self.logical_source.is_some() => { + SummaryExpr::BinaryOp { + lhs, + rhs, + operator, + timing: planner_types::post_asap::ExecutionTiming::ReadTime, + } if self.logical_source.is_some() => { let operator = logical::binary_operator(operator)?; QueryPlanNode::Logical { operator, @@ -680,7 +685,12 @@ where inputs: vec![self.lower(child)?], } } - SummaryExpr::BinaryOp { lhs, rhs, operator } if exact_value_executable(node) => { + SummaryExpr::BinaryOp { + lhs, + rhs, + operator, + timing: planner_types::post_asap::ExecutionTiming::ReadTime, + } if exact_value_executable(node) => { let planner_types::pre_asap::BinaryOpKind::Arithmetic(operator) = &operator.kind else { unreachable!() @@ -931,7 +941,12 @@ pub(crate) fn exact_value_executable(node: &SummaryNode) -> bool { exact_accumulator_value_source(node).is_some_and(exact_value_executable) } SummaryExpr::KeepPreAsap(expr) => scalar_literal(expr).is_some(), - SummaryExpr::BinaryOp { lhs, rhs, operator } => { + SummaryExpr::BinaryOp { + lhs, + rhs, + operator, + timing: planner_types::post_asap::ExecutionTiming::ReadTime, + } => { matches!( operator.kind, planner_types::pre_asap::BinaryOpKind::Arithmetic(_) diff --git a/control_plane/tests/offline_evidence.rs b/control_plane/tests/offline_evidence.rs index b36bd3fdb..0c4a936b1 100644 --- a/control_plane/tests/offline_evidence.rs +++ b/control_plane/tests/offline_evidence.rs @@ -351,6 +351,7 @@ fn binary_summary_has_explicit_warm_tier_fallback() { let child = bound(&model()); let root = std::rc::Rc::new(SummaryNode { expr: SummaryExpr::BinaryOp { + timing: planner_types::post_asap::ExecutionTiming::ReadTime, lhs: child.clone(), rhs: child.clone(), operator: BinaryOperator { diff --git a/crates/asap_types/Cargo.toml b/crates/asap_types/Cargo.toml index b7f65f90b..b55a8e9ec 100644 --- a/crates/asap_types/Cargo.toml +++ b/crates/asap_types/Cargo.toml @@ -34,4 +34,4 @@ sha2 = "0.10" # exactly (`control_plane/Cargo.toml`) -- two different revs of the same # git dependency in one workspace resolve to two distinct Rust types that # won't unify. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "e78572c57e6c3eea2a9017f968fbc65a9196d56e" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21" } diff --git a/data_plane/Cargo.toml b/data_plane/Cargo.toml index bc967e4fb..9c9b1de0d 100644 --- a/data_plane/Cargo.toml +++ b/data_plane/Cargo.toml @@ -39,8 +39,8 @@ sha2 = "0.10" # reduction: Reduction, .. }`) are `pre_asap` types, in the same crate now # (not a separate `asap-ir` import). Query serving consumes the compiled # QueryPlan; these types are used at physical-plan compilation boundaries. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "e78572c57e6c3eea2a9017f968fbc65a9196d56e" } -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "e78572c57e6c3eea2a9017f968fbc65a9196d56e" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21" } # Shared external (workspace) serde.workspace = true @@ -133,7 +133,7 @@ fs2 = "0.4" # none of them. [dev-dependencies] -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "e78572c57e6c3eea2a9017f968fbc65a9196d56e" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21" } tempfile = "3.20.0" criterion = { version = "0.5", features = ["html_reports"] } tokio-tungstenite = "0.21" diff --git a/data_plane/src/precompute_engine/maintenance_runtime.rs b/data_plane/src/precompute_engine/maintenance_runtime.rs index 0cfbcdb78..2b322ef25 100644 --- a/data_plane/src/precompute_engine/maintenance_runtime.rs +++ b/data_plane/src/precompute_engine/maintenance_runtime.rs @@ -29,6 +29,7 @@ enum MaintenanceValue { Rows { values: Vec<(i64, f64)>, name: String, + timestamped: bool, }, } @@ -125,6 +126,19 @@ impl PrecomputeOperatorRegistry for OperatorAdapter<'_> { ) -> Result { match &node.payload { ExecutableOperatorPayload::SummaryMerge => merge_inputs(inputs), + ExecutableOperatorPayload::Binary { + operator, + timing: planner_types::post_asap::ExecutionTiming::MaintenanceTime, + } => { + if !matches!(&self.inputs, MaintenanceInputs::Frozen(_)) + || node.output_state + != planner_types::post_asap::ExecutionDataState::MAINTENANCE_ROWS + { + return Err("maintenance binary requires immutable completed row inputs".into()); + } + evaluate_aligned_binary(node, operator, inputs) + } + ExecutableOperatorPayload::Value { operation: planner_types::post_asap::ValueOperation::FinalizeExactAccumulator, timing: planner_types::post_asap::ExecutionTiming::MaintenanceTime, @@ -141,7 +155,7 @@ impl PrecomputeOperatorRegistry for OperatorAdapter<'_> { let [value] = inputs else { return Err("maintenance SummaryAgg requires exactly one row input".into()); }; - let MaintenanceValue::Rows { values, name } = value.as_ref() else { + let MaintenanceValue::Rows { values, name, .. } = value.as_ref() else { return Err("maintenance SummaryAgg requires a typed update evaluator; finalize summary state before applying an update".into()); }; if !matches!(&self.inputs, MaintenanceInputs::Frozen(_)) { @@ -263,6 +277,119 @@ fn evaluate_weight( Ok(weight) } +// Rows carry only a value and its window timestamp. Reject any schema that +// would require silently dropping another column or manufacturing a timestamp. +fn maintenance_float64_column( + node: &ExecutableDagNode, +) -> Result<&planner_types::post_asap::SummaryField, String> { + use planner_types::post_asap::SummaryFamilyType; + let fields = &node.output_schema.fields; + let field = match node.output_schema.time_index { + None if fields.len() == 1 => &fields[0], + Some(time_index) if fields.len() == 2 && time_index < 2 => { + let timestamp = &fields[time_index]; + let value = &fields[1 - time_index]; + if timestamp.nullable + || timestamp.name == value.name + || !matches!( + timestamp.dtype, + SummaryFamilyType::Plain(planner_types::pre_asap::DataType::Timestamp) + ) + { + return Err("finalization timestamp column differs from its typed schema".into()); + } + value + } + _ => { + return Err( + "exact maintenance finalization requires one value and optional declared timestamp" + .into(), + ) + } + }; + if field.nullable + || !matches!( + field.dtype, + SummaryFamilyType::Plain(planner_types::pre_asap::DataType::Float64) + ) + { + return Err( + "exact maintenance finalization currently requires a Float64 output column".into(), + ); + } + Ok(field) +} + +fn evaluate_aligned_binary( + node: &ExecutableDagNode, + operator: &planner_types::post_asap::BinaryOperator, + inputs: &[Arc], +) -> Result { + use planner_types::pre_asap::BinaryOpKind; + let BinaryOpKind::Arithmetic(arithmetic) = &operator.kind else { + return Err("maintenance binary currently requires arithmetic".into()); + }; + if operator.vector_match.is_some() { + return Err( + "maintenance binary requires explicit population routing for vector matching".into(), + ); + } + let name = maintenance_float64_column(node)?.name.clone(); + if node.output_schema.time_index.is_none() { + return Err("maintenance binary requires declared window timestamps".into()); + } + let [left, right] = inputs else { + return Err("maintenance binary requires two row inputs".into()); + }; + let ( + MaintenanceValue::Rows { + values: left, + timestamped: true, + .. + }, + MaintenanceValue::Rows { + values: right, + timestamped: true, + .. + }, + ) = (left.as_ref(), right.as_ref()) + else { + return Err("maintenance binary requires finalized row inputs".into()); + }; + if left.is_empty() || left.len() != right.len() { + return Err("maintenance binary requires matching nonempty timestamp sets".into()); + } + // Canonical timestamp maps accept arrival-order differences, but never + // collapse duplicate updates or pair unrelated source windows by position. + let mut left_rows = BTreeMap::new(); + let mut right_rows = BTreeMap::new(); + for (rows, index) in [(left, &mut left_rows), (right, &mut right_rows)] { + for &(timestamp, value) in rows { + if !value.is_finite() || index.insert(timestamp, value).is_some() { + return Err( + "maintenance binary input has duplicate timestamps or non-finite values".into(), + ); + } + } + } + let mut values = Vec::with_capacity(left.len()); + for (timestamp, left) in left_rows { + let right = right_rows + .get(×tamp) + .ok_or("maintenance binary requires matching timestamp sets")?; + let value = crate::utils::arithmetic::evaluate_float64_arithmetic(arithmetic, left, *right); + if !value.is_finite() { + return Err("maintenance binary produced a non-finite update".into()); + } + values.push((timestamp, value)); + } + Ok(MaintenanceValue::Rows { + values, + name, + timestamped: true, + }) +} + fn finalize_exact( node: &ExecutableDagNode, inputs: &[Arc], @@ -303,40 +430,7 @@ fn finalize_exact( ) } }; - let fields = &node.output_schema.fields; - let field = match node.output_schema.time_index { - None if fields.len() == 1 => &fields[0], - Some(time_index) if fields.len() == 2 && time_index < 2 => { - let timestamp = &fields[time_index]; - let value = &fields[1 - time_index]; - if timestamp.nullable - || timestamp.name == value.name - || !matches!( - timestamp.dtype, - SummaryFamilyType::Plain(planner_types::pre_asap::DataType::Timestamp) - ) - { - return Err("finalization timestamp column differs from its typed schema".into()); - } - value - } - _ => { - return Err( - "exact maintenance finalization requires one value and optional declared timestamp" - .into(), - ) - } - }; - if field.nullable - || !matches!( - field.dtype, - SummaryFamilyType::Plain(planner_types::pre_asap::DataType::Float64) - ) - { - return Err( - "exact maintenance finalization currently requires a Float64 output column".into(), - ); - } + let field = maintenance_float64_column(node)?; let values = states .into_iter() .map(|(timestamp, state)| { @@ -352,6 +446,8 @@ fn finalize_exact( Ok(MaintenanceValue::Rows { values, name: field.name.clone(), + timestamped: node.output_schema.time_index.is_some() + && matches!(input.as_ref(), MaintenanceValue::SummaryWindows { .. }), }) } @@ -2416,6 +2512,110 @@ mod tests { persistence.shutdown(); } + #[test] + fn maintenance_binary_aligns_windows_and_rejects_incomplete_or_ambiguous_rows() { + use planner_types::post_asap::{BinaryOperator, SummaryFamilyType, SummaryField}; + use planner_types::pre_asap::{ArithmeticOpKind, BinaryOpKind, DataType}; + let mut operation = node(10); + operation.output_schema.fields = vec![ + SummaryField { + name: "ts".into(), + dtype: SummaryFamilyType::Plain(DataType::Timestamp), + nullable: false, + }, + SummaryField { + name: "value".into(), + dtype: SummaryFamilyType::Plain(DataType::Float64), + nullable: false, + }, + ]; + operation.output_schema.time_index = Some(0); + let mut operator = BinaryOperator { + kind: BinaryOpKind::Arithmetic(ArithmeticOpKind::Sub), + vector_match: None, + }; + let rows = |values: Vec<(i64, f64)>| { + Arc::new(MaintenanceValue::Rows { + values, + name: "value".into(), + timestamped: true, + }) + }; + let left = rows(vec![(2_000, 7.0), (1_000, 5.0)]); + let right = rows(vec![(1_000, 2.0), (2_000, 3.0)]); + // Arrival order cannot exchange windows, and subtraction retains edge order. + let MaintenanceValue::Rows { values, .. } = + evaluate_aligned_binary(&operation, &operator, &[left.clone(), right.clone()]).unwrap() + else { + panic!("expected rows") + }; + assert_eq!(values, vec![(1_000, 3.0), (2_000, 4.0)]); + let binding = BackendExecutableBinding { + nodes: BTreeMap::new(), + query_sink: PostAsapNodeId(10), + query_plan_sink: asap_types::query_plan::QueryNodeId(10), + precompute_sinks: vec![], + }; + let frozen = OperatorAdapter { + binding: &binding, + inputs: MaintenanceInputs::Frozen(&[]), + configs: &[], + }; + operation.payload = ExecutableOperatorPayload::Binary { + operator: operator.clone(), + timing: planner_types::post_asap::ExecutionTiming::MaintenanceTime, + }; + operation.output_state = planner_types::post_asap::ExecutionDataState::MAINTENANCE_ROWS; + assert!(frozen + .execute(&operation, &[left.clone(), right.clone()]) + .is_ok()); + let live = OperatorAdapter { + binding: &binding, + inputs: MaintenanceInputs::Live { + definition: definition(1), + state: sum(1.0), + }, + configs: &[], + }; + assert!(live + .execute(&operation, &[left.clone(), right.clone()]) + .is_err()); + operation.payload = ExecutableOperatorPayload::Binary { + operator: operator.clone(), + timing: planner_types::post_asap::ExecutionTiming::ReadTime, + }; + assert!(frozen + .execute(&operation, &[left.clone(), right.clone()]) + .is_err()); + + for invalid in [ + rows(vec![]), + rows(vec![(1_000, 2.0)]), + rows(vec![(1_000, 2.0), (3_000, 3.0)]), + rows(vec![(1_000, 2.0), (1_000, 3.0)]), + rows(vec![(1_000, f64::NAN), (2_000, 3.0)]), + Arc::new(MaintenanceValue::Rows { + values: vec![(1_000, 2.0), (2_000, 3.0)], + name: "value".into(), + timestamped: false, + }), + Arc::new(MaintenanceValue::summary(sum(2.0))), + ] { + assert!( + evaluate_aligned_binary(&operation, &operator, &[left.clone(), invalid]).is_err() + ); + } + operator.kind = BinaryOpKind::Arithmetic(ArithmeticOpKind::Div); + assert!(evaluate_aligned_binary( + &operation, + &operator, + &[left.clone(), rows(vec![(1_000, 0.0), (2_000, 3.0)])] + ) + .is_err()); + operation.output_schema.fields[1].dtype = SummaryFamilyType::Plain(DataType::Int64); + assert!(evaluate_aligned_binary(&operation, &operator, &[left, right]).is_err()); + } + #[test] fn finalization_preserves_windows_until_an_explicit_merge() { use planner_types::post_asap::{ExactKind, ExactParams, SummaryFamilyType, SummaryField}; @@ -2480,7 +2680,7 @@ mod tests { }, ]; read.output_schema.time_index = Some(0); - let MaintenanceValue::Rows { values, name } = + let MaintenanceValue::Rows { values, name, .. } = finalize_exact(&read, &[Arc::clone(&input)]).unwrap() else { panic!("expected typed rows") diff --git a/data_plane/src/precompute_engine/subdag_scheduler.rs b/data_plane/src/precompute_engine/subdag_scheduler.rs index 6a52b88df..69d67c401 100644 --- a/data_plane/src/precompute_engine/subdag_scheduler.rs +++ b/data_plane/src/precompute_engine/subdag_scheduler.rs @@ -1,6 +1,8 @@ use asap_types::executable_plan::{BackendExecutableBinding, BackendNodeBinding}; use planner_types::post_asap::PostAsapNodeId; -use planner_types::post_asap::{ExecutableDag, ExecutableDagNode, ExecutionDataState}; +use planner_types::post_asap::{ + EdgeRole, ExecutableDag, ExecutableDagNode, ExecutableOperatorPayload, ExecutionDataState, +}; use std::{ collections::{BTreeMap, BTreeSet}, sync::Arc, @@ -91,12 +93,41 @@ where .iter() .map(|n| (n.id.0, n)) .collect::>(); - let mut inputs = BTreeMap::>::new(); + let mut incoming = BTreeMap::>::new(); for edge in &dag.edges { - inputs - .entry(edge.consumer.0) - .or_default() - .push(edge.producer.0); + incoming.entry(edge.consumer.0).or_default().push(edge); + } + for node in &dag.nodes { + if matches!(node.payload, ExecutableOperatorPayload::Binary { .. }) { + incoming.entry(node.id.0).or_default(); + } + } + let mut inputs = BTreeMap::>::new(); + for (consumer, edges) in incoming { + let ordered = if nodes + .get(&consumer) + .is_some_and(|node| matches!(node.payload, ExecutableOperatorPayload::Binary { .. })) + { + // Wire order is not operand order. Preserve noncommutative binary + // semantics even when a valid transport reorders its edge list. + let left = edges + .iter() + .filter(|edge| edge.role == EdgeRole::Left) + .collect::>(); + let right = edges + .iter() + .filter(|edge| edge.role == EdgeRole::Right) + .collect::>(); + if edges.len() != 2 || left.len() != 1 || right.len() != 1 { + return Err(ScheduleError::Invalid( + "binary maintenance input roles must be exactly Left and Right".into(), + )); + } + vec![left[0].producer.0, right[0].producer.0] + } else { + edges.iter().map(|edge| edge.producer.0).collect() + }; + inputs.insert(consumer, ordered); } let mut active = BTreeSet::new(); let mut values = BTreeMap::>::new(); @@ -272,6 +303,69 @@ mod tests { } } + #[test] + fn binary_operand_roles_survive_edge_reordering_and_reject_duplicates() { + use planner_types::post_asap::BinaryOperator; + use planner_types::pre_asap::{ArithmeticOpKind, BinaryOpKind}; + struct Subtract; + impl PrecomputeOperatorRegistry for Subtract { + type Error = String; + fn execute( + &self, + node: &ExecutableDagNode, + inputs: &[Arc], + ) -> Result { + match node.id.0 { + 0 => Ok(10), + 1 => Ok(3), + 3 => Ok(*inputs[0] - *inputs[1]), + _ => Err("unexpected node".into()), + } + } + } + let mut binary = node(3); + binary.operator = ExecutableOperator::Binary; + binary.payload = ExecutableOperatorPayload::Binary { + timing: planner_types::post_asap::ExecutionTiming::MaintenanceTime, + operator: BinaryOperator { + kind: BinaryOpKind::Arithmetic(ArithmeticOpKind::Sub), + vector_match: None, + }, + }; + let mut left = edge(0, 3); + left.role = EdgeRole::Left; + let mut right = edge(1, 3); + right.role = EdgeRole::Right; + let mut dag = ExecutableDag { + nodes: vec![node(0), node(1), node(2), binary, { + let mut query = node(4); + query.output_state = ExecutionDataState::READ_ROWS; + query + }], + edges: vec![right, left], + root: PostAsapNodeId(3), + }; + let execute = |dag: &ExecutableDag| { + execute_precompute_sink( + dag, + &binding(), + PostAsapNodeId(3), + key(3), + &Subtract, + &Sink::default(), + ) + }; + assert_eq!(*execute(&dag).unwrap(), 7); + dag.edges.reverse(); + assert_eq!(*execute(&dag).unwrap(), 7); + dag.edges[1].role = EdgeRole::Left; + assert!(matches!(execute(&dag), Err(ScheduleError::Invalid(_)))); + dag.edges.pop(); + assert!(matches!(execute(&dag), Err(ScheduleError::Invalid(_)))); + dag.edges.clear(); + assert!(matches!(execute(&dag), Err(ScheduleError::Invalid(_)))); + } + #[test] fn shared_dependency_executes_once_and_replay_reads_committed_value() { // 0 is shared by 1 and 2; 3 consumes both branches. diff --git a/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs b/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs index 8e1184736..6a4c62683 100644 --- a/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs +++ b/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs @@ -4,6 +4,8 @@ use std::collections::BTreeMap; +use crate::utils::arithmetic::evaluate_float64_arithmetic as arithmetic; + use crate::query_engines::asap_query_engine::summary_exec::{execute, ExecOutcome}; use asap_types::query_plan::{QueryNodeId, QueryPlanNode}; use control_plane::types_v2::AccuracyTarget; @@ -410,19 +412,6 @@ fn expand_item_readout( Ok((expand_item_rows(group_key, value, item_labels)?, coverage)) } -fn arithmetic(operator: &planner_types::pre_asap::ArithmeticOpKind, left: f64, right: f64) -> f64 { - use planner_types::pre_asap::ArithmeticOpKind::*; - match operator { - Add => left + right, - Sub => left - right, - Mul => left * right, - Div => left / right, - Mod => left % right, - Pow => left.powf(right), - Atan2 => left.atan2(right), - } -} - fn binary_values( operator: &planner_types::pre_asap::ArithmeticOpKind, lhs: &PhysicalQueryOutput, diff --git a/data_plane/src/query_engines/asap_query_engine/summary_exec.rs b/data_plane/src/query_engines/asap_query_engine/summary_exec.rs index 4ab734ea4..672363d12 100644 --- a/data_plane/src/query_engines/asap_query_engine/summary_exec.rs +++ b/data_plane/src/query_engines/asap_query_engine/summary_exec.rs @@ -544,6 +544,7 @@ mod tests { let child = logical_node(); let tree = SummaryNode { expr: SummaryExpr::BinaryOp { + timing: planner_types::post_asap::ExecutionTiming::ReadTime, lhs: child.clone(), rhs: child.clone(), operator: planner_types::post_asap::BinaryOperator { diff --git a/data_plane/src/utils/arithmetic.rs b/data_plane/src/utils/arithmetic.rs new file mode 100644 index 000000000..94aba683f --- /dev/null +++ b/data_plane/src/utils/arithmetic.rs @@ -0,0 +1,19 @@ +//! Float64 arithmetic shared by data-plane execution engines. +//! Preserve IEEE non-finite results; callers own their output policies. + +pub(crate) fn evaluate_float64_arithmetic( + operator: &planner_types::pre_asap::ArithmeticOpKind, + left: f64, + right: f64, +) -> f64 { + use planner_types::pre_asap::ArithmeticOpKind::*; + match operator { + Add => left + right, + Sub => left - right, + Mul => left * right, + Div => left / right, + Mod => left % right, + Pow => left.powf(right), + Atan2 => left.atan2(right), + } +} diff --git a/data_plane/src/utils/mod.rs b/data_plane/src/utils/mod.rs index 5d620636b..72f331c5d 100644 --- a/data_plane/src/utils/mod.rs +++ b/data_plane/src/utils/mod.rs @@ -1,3 +1,4 @@ +pub(crate) mod arithmetic; pub mod file_io; pub mod http;