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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
52 changes: 52 additions & 0 deletions crates/asap-aware-mapping/src/replacement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1796,6 +1796,7 @@ fn realize_binary(

Ok(Some(Rc::new(SummaryNode {
expr: SummaryExpr::BinaryOp {
timing: ExecutionTiming::ReadTime,
lhs: lhs_node,
rhs: rhs_node,
operator: asap_types::post_asap::BinaryOperator {
Expand Down Expand Up @@ -2178,6 +2179,56 @@ fn realize_physical_summary_input(

/// Emit `SummaryAgg` (recursively binding the child), plus the
/// `SummaryEstimate` readout when `estimate` is set.
// Retain the exact expression and schema while placing its value production
// on the update path. Read-time consumers keep their original shared nodes.
fn maintenance_exact_values(node: Rc<SummaryNode>) -> Rc<SummaryNode> {
let expr = match &node.expr {
SummaryExpr::BinaryOp {
lhs, rhs, operator, ..
} if operator.vector_match.is_none()
&& matches!(
operator.kind,
asap_types::pre_asap::BinaryOpKind::Arithmetic(_)
)
&& node
.guarantee
.as_ref()
.is_some_and(ResultGuarantee::is_exact) =>
{
SummaryExpr::BinaryOp {
lhs: maintenance_exact_values(lhs.clone()),
rhs: maintenance_exact_values(rhs.clone()),
operator: operator.clone(),
timing: ExecutionTiming::MaintenanceTime,
}
}
SummaryExpr::ValueOperation {
child,
operation: ValueOperation::FinalizeExactAccumulator,
..
} if matches!(
child.expr,
SummaryExpr::SummaryAgg {
family: SummaryFamilyType::ExactAggregate(..),
..
}
) =>
{
SummaryExpr::ValueOperation {
child: child.clone(),
operation: ValueOperation::FinalizeExactAccumulator,
timing: ExecutionTiming::MaintenanceTime,
}
}
_ => return node,
};
Rc::new(SummaryNode {
expr,
schema: node.schema.clone(),
guarantee: node.guarantee.clone(),
})
}

#[allow(clippy::too_many_arguments)]
fn construct_summary_agg(
node: &QueryExpr,
Expand Down Expand Up @@ -2241,6 +2292,7 @@ fn construct_summary_agg(
// an exact scalar accumulator currently stores its value directly.
let bound_child =
finalize_exact_accumulator_at(bound_child, &input.child, ExecutionTiming::MaintenanceTime)?;
let bound_child = maintenance_exact_values(bound_child);

// ── Guarantee (issue #172) ──────────────────────────────────────────
// Derived *before* the node exists, so an illegal composition is never
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3000,6 +3000,7 @@ mod tests {
let operand = summary_with_operations(false, false, false);
Rc::new(SummaryNode {
expr: SummaryExpr::BinaryOp {
timing: asap_types::post_asap::ExecutionTiming::ReadTime,
lhs: Rc::clone(&operand),
rhs: operand,
operator: asap_types::post_asap::BinaryOperator {
Expand Down
38 changes: 38 additions & 0 deletions crates/integration-tests/tests/promql_to_post_asap.rs
Original file line number Diff line number Diff line change
Expand Up @@ -943,3 +943,41 @@ fn nested_summary_explicitly_finalizes_exact_child_at_maintenance_time() {
.any(|field| matches!(field.dtype, SummaryFamilyType::Plain(DataType::Float64))));
compile_executable_dag(&plan).expect("explicit boundary is a valid executable DAG");
}

#[test]
fn exact_binary_maintenance_has_explicit_timing_and_legacy_wire_default() {
use asap_types::post_asap::{ExecutableOperatorPayload, ExecutionTiming};
for (query, expected) in [
(
"quantile(0.9, sum_over_time(m[1m]) + sum_over_time(n[1m]))",
ExecutionTiming::MaintenanceTime,
),
(
"sum_over_time(m[1m]) + sum_over_time(n[1m])",
ExecutionTiming::ReadTime,
),
] {
let input = lower_promql(query, AccuracyTarget::Epsilon(0.05)).unwrap();
let search = search_workload(vec![("q", Rc::new(input))]);
let choice = search.global_selection(&DefaultCostModel);
let plan = choice.materialize(&search.roots[0].1).unwrap().unwrap();
let dag = compile_executable_dag(&plan).unwrap();
let payload = dag
.nodes
.iter()
.find_map(|node| {
matches!(node.payload, ExecutableOperatorPayload::Binary { .. })
.then_some(&node.payload)
})
.unwrap();
assert!(
matches!(payload, ExecutableOperatorPayload::Binary { timing, .. } if *timing == expected)
);
let wire = serde_json::to_value(payload).unwrap();
if expected == ExecutionTiming::ReadTime {
assert!(wire.get("timing").is_none());
}
let restored: ExecutableOperatorPayload = serde_json::from_value(wire).unwrap();
assert_eq!(&restored, payload);
}
}
5 changes: 4 additions & 1 deletion crates/types/src/post_asap/cse.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,13 +31,15 @@ fn same_node(left: &SummaryNode, right: &SummaryNode) -> bool {
lhs: al,
rhs: ar,
operator: ao,
timing: at,
},
BinaryOp {
lhs: bl,
rhs: br,
operator: bo,
timing: bt,
},
) => Rc::ptr_eq(al, bl) && Rc::ptr_eq(ar, br) && ao == bo,
) => Rc::ptr_eq(al, bl) && Rc::ptr_eq(ar, br) && ao == bo && at == bt,
(
CandidateTopK {
candidates: ac,
Expand Down Expand Up @@ -411,6 +413,7 @@ mod tests {
for _ in 0..24 {
current = Rc::new(SummaryNode {
expr: SummaryExpr::BinaryOp {
timing: super::super::ExecutionTiming::ReadTime,
lhs: Rc::clone(&current),
rhs: current,
operator: super::super::BinaryOperator {
Expand Down
7 changes: 6 additions & 1 deletion crates/types/src/post_asap/executable_dag.rs
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,8 @@ pub enum ExecutableOperatorPayload {
expression: QueryExpr,
},
Binary {
#[serde(default, skip_serializing_if = "ExecutionTiming::is_read_time")]
timing: ExecutionTiming,
operator: BinaryOperator,
},
CandidateTopK {
Expand Down Expand Up @@ -444,7 +446,10 @@ pub fn compile_executable_dag_with_node_ids(
SummaryExpr::KeepPreAsap(expression) => ExecutableOperatorPayload::Fallback {
expression: (**expression).clone(),
},
SummaryExpr::BinaryOp { operator, .. } => ExecutableOperatorPayload::Binary {
SummaryExpr::BinaryOp {
operator, timing, ..
} => ExecutableOperatorPayload::Binary {
timing: *timing,
operator: operator.clone(),
},
SummaryExpr::CandidateTopK {
Expand Down
71 changes: 62 additions & 9 deletions crates/types/src/post_asap/execution_data_state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,13 +52,19 @@ use crate::pre_asap::query_expr::{aggregate_output_schema, QueryExprError};
use crate::pre_asap::schema::{Column, Schema};

/// When a post-ASAP value is produced.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
#[derive(
Default, Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize,
)]
pub enum ExecutionTiming {
MaintenanceTime,
#[default]
ReadTime,
}

impl ExecutionTiming {
pub fn is_read_time(&self) -> bool {
*self == Self::ReadTime
}
pub fn as_str(self) -> &'static str {
match self {
Self::MaintenanceTime => "maintenance_time",
Expand Down Expand Up @@ -188,6 +194,8 @@ pub enum ExecutionDataStateError {
/// nothing maintains state above it, so its output is never read.
#[error("A maintenance-time value operation cannot be a plan root: its update-path output feeds nothing")]
MaintenanceRowsAtRoot,
#[error("unsupported maintenance binary schema or operator")]
InvalidMaintenanceBinary,
/// An `ExactOperation` whose input columns are not all `Plain` at its
/// declared data_state.
#[error("exact operator consumes non-plain column {column:?} ({dtype})")]
Expand Down Expand Up @@ -223,9 +231,13 @@ impl ExecutionDataStateAssignment {
pub fn produced_data_state(expr: &SummaryExpr) -> Option<ExecutionDataState> {
Some(match expr {
SummaryExpr::KeepPreAsap(_) => return None,
SummaryExpr::BinaryOp { .. }
| SummaryExpr::CandidateTopK { .. }
| SummaryExpr::RelationalJoin { .. } => ExecutionDataState::READ_ROWS,
SummaryExpr::BinaryOp { timing, .. } => ExecutionDataState {
timing: *timing,
primitive: DataPrimitive::Raw,
},
SummaryExpr::CandidateTopK { .. } | SummaryExpr::RelationalJoin { .. } => {
ExecutionDataState::READ_ROWS
}
SummaryExpr::SummaryAgg { .. }
| SummaryExpr::SummaryJoin { .. }
| SummaryExpr::SummarySubtract { .. }
Expand Down Expand Up @@ -312,11 +324,45 @@ fn visit(

match &node.expr {
SummaryExpr::KeepPreAsap(_) => Ok(()),
SummaryExpr::BinaryOp { lhs, rhs, .. } => {
SummaryExpr::BinaryOp {
lhs,
rhs,
timing,
operator,
} => {
if *timing == ExecutionTiming::MaintenanceTime {
use crate::pre_asap::{BinaryOpKind, DataType};
if operator.vector_match.is_some()
|| !matches!(operator.kind, BinaryOpKind::Arithmetic(_))
|| lhs.schema != rhs.schema
|| lhs.schema != node.schema
|| !node.schema.fields.iter().all(|field| {
!field.nullable
&& matches!(
field.dtype,
SummaryFamilyType::Plain(DataType::Float64 | DataType::Timestamp)
)
})
|| node
.schema
.fields
.iter()
.filter(|field| {
matches!(field.dtype, SummaryFamilyType::Plain(DataType::Float64))
})
.count()
!= 1
{
return Err(ExecutionDataStateError::InvalidMaintenanceBinary);
}
}
let expected = ExecutionDataState {
timing: *timing,
primitive: DataPrimitive::Raw,
};
for input in [lhs, rhs] {
let state =
produced_data_state(&input.expr).unwrap_or(ExecutionDataState::READ_ROWS);
if state != ExecutionDataState::READ_ROWS {
let state = produced_data_state(&input.expr).unwrap_or(expected);
if state != expected {
return Err(ExecutionDataStateError::IllegalChildDataState {
edge: "BinaryOp operand",
child: state,
Expand Down Expand Up @@ -453,9 +499,16 @@ pub fn assigned_child_data_state(parent: &SummaryExpr, child: &SummaryNode) -> E
SummaryExpr::ValueOperation {
timing: ExecutionTiming::ReadTime,
..
}
| SummaryExpr::BinaryOp {
timing: ExecutionTiming::ReadTime,
..
} => ExecutionDataState::READ_ROWS,
SummaryExpr::KeepPreAsap(_)
| SummaryExpr::BinaryOp { .. }
| SummaryExpr::BinaryOp {
timing: ExecutionTiming::MaintenanceTime,
..
}
| SummaryExpr::CandidateTopK { .. }
| SummaryExpr::RelationalJoin { .. }
| SummaryExpr::SummaryAgg { .. }
Expand Down
2 changes: 2 additions & 0 deletions crates/types/src/post_asap/expr.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
use super::ExecutionTiming;
use std::rc::Rc;

use super::guarantee::ResultGuarantee;
Expand Down Expand Up @@ -120,6 +121,7 @@ pub enum SummaryExpr {
/// This keeps realizable summary/readout leaves visible instead of
/// hiding the complete expression inside `KeepPreAsap`.
BinaryOp {
timing: ExecutionTiming,
lhs: Rc<SummaryNode>,
rhs: Rc<SummaryNode>,
operator: BinaryOperator,
Expand Down
Loading