Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
48 commits
Select commit Hold shift + click to select a range
215fc79
refactor: execute query relation subgraphs with memoization
zzylol Sep 22, 2026
e2df4a5
docs: explain installed QueryPlan DAG execution
zzylol Sep 22, 2026
f9bb469
chore: sync backend with ASAPPlanner main
zzylol Sep 22, 2026
7435dc5
feat: cover post-ASAP query operators exhaustively
zzylol Sep 22, 2026
aca4349
fix: leave Planner evidence upgrade to stacked PR
zzylol Sep 22, 2026
e1fa604
fix: preserve incompatible Planner sync PRs
zzylol Sep 22, 2026
78d2026
Merge branch 'stack/requested-763' into stack/requested-765
zzylol Sep 22, 2026
18fbbda
Merge branch 'stack/requested-763' into stack/requested-765
zzylol Sep 22, 2026
9a362ec
Merge branch 'stack/planner-api-763' into stack/planner-api-765
zzylol Sep 23, 2026
5781cfb
Merge branch 'stack/planner-api-763' into stack/planner-api-765
zzylol Sep 23, 2026
63ec16f
Merge branch 'stack/planner-api-763' into stack/planner-api-765
zzylol Sep 23, 2026
2a2257a
Merge branch 'stack/planner-api-763' into stack/planner-api-765
zzylol Sep 23, 2026
e1ebed0
Merge branch 'stack/planner-api-763' into stack/planner-api-765
zzylol Sep 23, 2026
6219ea5
Merge branch 'stack/planner-api-763' into stack/planner-api-765
zzylol Sep 23, 2026
1edde79
refactor: share physical operator kernels across ASAP deployments
zzylol Sep 23, 2026
1a3af01
test: include shared physical kernels in the contracts suite
zzylol Sep 23, 2026
dccbe18
fix: validate native kernel dimensions independently of packed codecs
zzylol Sep 23, 2026
c550822
docs: specify physical operator coverage and execution gaps
zzylol Sep 23, 2026
935ffe0
docs: clarify Planner execution phase terminology
zzylol Sep 23, 2026
f8fbaaa
refactor!: compose membership filtering and TopK with ingestion and q…
zzylol Sep 23, 2026
e0bb3d3
test: specify ingestion phase in derived merge fixtures
zzylol Sep 23, 2026
cfee4df
fix: exclude query-produced states from ingestion scheduling
zzylol Sep 23, 2026
ab6a999
docs: explain remaining query-time computations in plain language
zzylol Sep 23, 2026
22a6908
docs: consolidate physical operation coverage and defer raw scan
zzylol Sep 23, 2026
309bca0
docs: define precomputation modes independently of KLL
zzylol Sep 23, 2026
0050c53
docs: align query execution diagram with physical operation terminology
zzylol Sep 23, 2026
f0084ad
docs: clarify current-series state operations
zzylol Sep 23, 2026
5aded53
docs: define independent DAG runtime shared by both engines
zzylol Sep 23, 2026
1591b28
Implement independent shared physical DAG runtime and native operators
zzylol Sep 23, 2026
17de4a5
Bind ingestion and query value computation to native DAG operators
zzylol Sep 23, 2026
fe840ef
Stack query DAG integration on native precompute and shared operators
zzylol Sep 23, 2026
3a0f588
Use native scalar sources and document engine operator bindings
zzylol Sep 23, 2026
cfde882
Merge branch 'stack/native-precompute-763' into refactor/native-engin…
zzylol Sep 23, 2026
39d74f7
Merge branch 'stack/native-precompute-763' into refactor/native-engin…
zzylol Sep 23, 2026
8710c86
Merge branch 'stack/native-precompute-763' into refactor/native-engin…
zzylol Sep 23, 2026
4b68111
Inherit formatted shared library and precompute boundaries
zzylol Sep 23, 2026
6a707f7
Merge branch 'stack/native-precompute-763' into refactor/native-engin…
zzylol Sep 23, 2026
2bf1a1d
Merge branch 'stack/native-precompute-763' into refactor/native-engin…
zzylol Sep 23, 2026
dc0b617
refactor: execute relation and temporal computations in Planner library
zzylol Sep 23, 2026
68283b8
Merge branch 'stack/native-precompute-763' into refactor/native-engin…
zzylol Sep 23, 2026
a13a603
Merge branch 'stack/native-precompute-763' into refactor/native-engin…
zzylol Sep 23, 2026
08f475b
fix: share storage frontier execution within query DAG scope
zzylol Sep 23, 2026
f54330c
Merge branch 'stack/native-precompute-763' into refactor/native-engin…
zzylol Sep 23, 2026
9830e24
Merge branch 'stack/native-precompute-763' into refactor/native-engin…
zzylol Sep 23, 2026
bba134e
fix: accept Planner integer score columns in query sorting
zzylol Sep 23, 2026
d341883
Merge branch 'stack/native-precompute-763' into refactor/native-engin…
zzylol Sep 23, 2026
72de955
Merge branch 'stack/native-precompute-763' into refactor/native-engin…
zzylol Sep 23, 2026
d5ec035
Merge branch 'restack/763' into restack/765
zzylol Sep 24, 2026
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
12 changes: 4 additions & 8 deletions control_plane/src/clickhouse.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1453,14 +1453,10 @@ mod tests {
("telemetry".into(), timestamped("timestamp_ms", "value")),
(
"divisors".into(),
Schema::with_time_index(
vec![
Column::new("timestamp", DataType::Int64, false),
Column::new("divisor", DataType::Float64, false),
],
0,
vec![],
),
Schema::new(vec![
Column::new("timestamp", DataType::Int64, false),
Column::new("divisor", DataType::Float64, false),
]),
),
]),
accuracy: AccuracyTarget::Exact,
Expand Down
45 changes: 9 additions & 36 deletions control_plane/src/physical/executable_binding.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,44 +4,17 @@ pub use asap_types::executable_plan::*;

#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum OperatorExecution {
Maintenance,
Ingestion,
Query,
}

/// Keep the backend's ownership decision exhaustive over Planner's physical IR.
/// Adding a payload variant upstream must therefore choose an executor here.
fn operator_execution(
node: &planner_types::post_asap::ExecutableDagNode,
) -> Result<OperatorExecution, String> {
use planner_types::post_asap::{ExecutableOperatorPayload as Payload, ExecutionTiming};

let declared = match &node.payload {
Payload::Binary { timing, .. }
| Payload::Value { timing, .. }
| Payload::SummaryMerge { timing } => *timing,
Payload::MembershipFilter { .. } | Payload::SummaryEstimate { .. } => {
ExecutionTiming::QueryTime
}
Payload::SummaryAgg { .. }
| Payload::SummaryJoin { .. }
| Payload::SummarySubtract
| Payload::SummaryDelete { .. } => ExecutionTiming::IngestionTime,
// These operators can be placed on either side of the stored-state
// boundary. Planner's validated output state is authoritative.
Payload::Fallback { .. } | Payload::RelationalJoin { .. } => node.output_state.timing,
};
if declared != node.output_state.timing {
return Err(format!(
"post-ASAP node {:?} has operator timing {} but output state {}",
node.id,
declared.as_str(),
node.output_state
));
/// Every physical operator uses its node's placement; payload kind does not
/// restrict execution phase. Runtime capability is checked separately.
fn operator_execution(node: &planner_types::post_asap::ExecutableDagNode) -> OperatorExecution {
match node.output_state.timing {
planner_types::post_asap::ExecutionTiming::IngestionTime => OperatorExecution::Ingestion,
planner_types::post_asap::ExecutionTiming::QueryTime => OperatorExecution::Query,
}
Ok(match declared {
ExecutionTiming::IngestionTime => OperatorExecution::Maintenance,
ExecutionTiming::QueryTime => OperatorExecution::Query,
})
}

/// Assign backend phases to a selected semantic DAG without changing its nodes.
Expand All @@ -67,11 +40,11 @@ pub fn install_selected_dag(
asap_physical_operators::capability::validate_summary_kernel(family, input, grouping)
.map_err(|reason| format!("post-ASAP node {:?}: {reason}", node.id))?;
}
let execution = operator_execution(node)?;
let execution = operator_execution(node);
let binding = if let Some(summary_definition) = materialization(node.id) {
precompute_sinks.push(node.id);
BackendNodeBinding::Materialization { summary_definition }
} else if execution == OperatorExecution::Maintenance {
} else if execution == OperatorExecution::Ingestion {
BackendNodeBinding::MaintenanceInput
} else {
query_node(node.id).map_or(BackendNodeBinding::QueryInput, |query_node| {
Expand Down
1 change: 1 addition & 0 deletions control_plane/src/query_plan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -386,6 +386,7 @@ where
child.schema.fields.get(column).map(|field| &field.dtype),
Some(SummaryFamilyType::Plain(
planner_types::pre_asap::DataType::Float64
| planner_types::pre_asap::DataType::Int64
)) | Some(SummaryFamilyType::ExactAggregate(..))
) {
return Err(QueryPlanError::Invalid(
Expand Down
97 changes: 97 additions & 0 deletions crates/asap_types/src/query_plan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -366,6 +366,7 @@ impl QueryPlanEntry {
)));
}
for (id, node) in &self.nodes {
validate_native_relation(*id, node)?;
if let QueryPlanNode::Logical { operator, inputs } = node {
operator.validate(inputs.len())?;
}
Expand Down Expand Up @@ -746,3 +747,99 @@ mod contract_tests {
assert_send_sync::<super::QueryPlanEntry>();
}
}

/// Bind portable relation semantics before an installed plan can access its sources.
fn validate_native_relation(id: QueryNodeId, node: &QueryPlanNode) -> Result<(), QueryPlanError> {
use planner_types::post_asap::{
ExecutableDagNode, ExecutableOperatorPayload as Payload, ExecutionDataState, PostAsapNodeId,
};
use std::sync::Arc;
let invalid = |error: String| QueryPlanError::Invalid(format!("query node {}: {error}", id.0));
let (payload, inputs, output) = match node {
QueryPlanNode::Relational {
operation,
input_schema,
output_schema,
..
} => (
Payload::Value {
operation: serde_json::from_value(operation.clone())
.map_err(|e| invalid(e.to_string()))?,
},
vec![Arc::new(input_schema.clone())],
output_schema,
),
QueryPlanNode::RelationalJoin {
join_kind,
pred,
left_schema,
right_schema,
output_schema,
pruning,
..
} => (
Payload::RelationalJoin {
join_kind: join_kind.clone(),
pred: serde_json::from_value(pred.clone()).map_err(|e| invalid(e.to_string()))?,
pruning: serde_json::from_value(
serde_json::to_value(pruning).map_err(|e| invalid(e.to_string()))?,
)
.map_err(|e| invalid(e.to_string()))?,
},
vec![
Arc::new(left_schema.clone()),
Arc::new(right_schema.clone()),
],
output_schema,
),
_ => return Ok(()),
};
let node = ExecutableDagNode {
id: PostAsapNodeId(0),
payload,
output_state: ExecutionDataState::QUERY_ROWS,
output_schema: output.clone(),
guarantee: None,
};
asap_physical_operators::dag::planner::bind_node(&node, &inputs)
.map_err(|e| invalid(e.to_string()))?;
Ok(())
}

#[cfg(test)]
mod native_binding_tests {
use super::*;
use planner_types::{
post_asap::{SummaryFamilyType, SummaryField, SummarySchema, ValueOperation},
pre_asap::{DataType, Predicate, QueryExpr},
};

// Unsupported expressions fail installation without evaluating any source.
#[test]
fn rejects_unimplemented_relation_predicate_before_execution() {
let schema = SummarySchema {
fields: vec![SummaryField {
name: "value".into(),
dtype: SummaryFamilyType::Plain(DataType::Float64),
nullable: false,
}],
time_index: None,
};
let node = QueryPlanNode::Relational {
input: QueryNodeId(0),
operation: serde_json::to_value(ValueOperation::Filter {
pred: Predicate(
QueryExpr::FunctionCall {
name: "unimplemented_predicate".into(),
args: vec![],
}
.into(),
),
})
.unwrap(),
input_schema: schema.clone(),
output_schema: schema,
};
assert!(validate_native_relation(QueryNodeId(1), &node).is_err());
}
}
6 changes: 3 additions & 3 deletions data_plane/src/precompute_engine/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4624,9 +4624,9 @@ mod dag_execution_tests {
.unwrap();
installed.document.schema_version =
asap_types::executable_plan::MAINTENANCE_DAG_SCHEMA_VERSION;
let error = InstalledPrecomputePlan::from_precompute_plan(plan)
assert!(InstalledPrecomputePlan::from_precompute_plan(plan)
.unwrap_err()
.to_string();
assert!(error.contains("update"), "{error}");
.to_string()
.contains("update"));
}
}
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
//! Catalog-backed ClickHouse acceleration boundary.

use asap_physical_operators::accumulators::SumAccumulator;
use async_trait::async_trait;
use axum::{
body::Bytes,
Expand Down Expand Up @@ -856,7 +857,12 @@ mod tests {
"SELECT sum(value) FROM requests WHERE timestamp >= 2000 AND timestamp < 3000".into();
let uncovered = accelerator.execute(&request).await;
assert!(
matches!(&uncovered, ClickHouseAccelerationOutcome::Fallback(ClickHouseAccelerationFallback::Execution(detail)) if detail.contains("NoCandidates")),
matches!(
&uncovered,
ClickHouseAccelerationOutcome::Fallback(
ClickHouseAccelerationFallback::IncompleteCoverage
)
),
"{uncovered:?}"
);
request.sql =
Expand Down
Loading
Loading