Skip to content
Open
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
41 changes: 27 additions & 14 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

15 changes: 9 additions & 6 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@
resolver = "2"
members = [
"crates/asap_otel_proto",
"crates/asap_sketch_codec",
"crates/asap_summary_state",
"crates/asap_types",
"data_plane",
"control_plane",
Expand All @@ -14,10 +16,10 @@ version = "0.1.0"
[workspace.dependencies]
# Keep Planner frontends, selection, and IR on the same immutable revision.
# Alias upstream asap-types because this workspace also defines asap_types.
planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "76fbbf16cc44b19f56a780bfdb47a95327e84711" }
asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "76fbbf16cc44b19f56a780bfdb47a95327e84711" }
asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "76fbbf16cc44b19f56a780bfdb47a95327e84711" }
asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "76fbbf16cc44b19f56a780bfdb47a95327e84711" }
planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c98281a59df59740f006ca0b5b8da1ea9bb7fc41" }
asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c98281a59df59740f006ca0b5b8da1ea9bb7fc41" }
asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c98281a59df59740f006ca0b5b8da1ea9bb7fc41" }
asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c98281a59df59740f006ca0b5b8da1ea9bb7fc41" }

# Shared external deps (used by 2+ crates)
serde = { version = "1.0", features = ["derive"] }
Expand All @@ -37,8 +39,9 @@ arc-swap = "1.7"
reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] }

# Internal crates
asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "76fbbf16cc44b19f56a780bfdb47a95327e84711" }
asap_sketch_codec = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "76fbbf16cc44b19f56a780bfdb47a95327e84711" }
asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c98281a59df59740f006ca0b5b8da1ea9bb7fc41" }
asap_sketch_codec = { path = "crates/asap_sketch_codec" }
asap_summary_state = { path = "crates/asap_summary_state" }
asap_types = { path = "crates/asap_types" }
asap_otel_proto = { path = "crates/asap_otel_proto" }
indexmap = { version = "2.0", features = ["serde"] }
2 changes: 1 addition & 1 deletion control_plane/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -199,7 +199,7 @@ struct CompileAndPublishPhysicalPlanRequest {
workload_cost_evidence: Option<physical::workload_cost::WorkloadCostEvidence>,
queries: Vec<PhysicalPlanQueryRequest>,
data_workload: planner_types::workload::DataWorkload,
dataset_identity: planner_types::post_asap::LogicalDatasetIdentity,
dataset_identity: asap_types::semantic_fragment::LogicalDatasetIdentity,
#[serde(rename = "collector_ids")]
target_collector_ids: Vec<String>,
capability_snapshot_id: String,
Expand Down
28 changes: 17 additions & 11 deletions control_plane/src/physical/compiler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
//! Planner owns semantic candidates and guarantees. This module owns the
//! deployment decision: evidence freshness, target capabilities, windows, the
//! Collector execution projection, SummaryCatalog, and executable plans.
use asap_types::physical_plan_codec::PhysicalPlanCodec;

use std::collections::{BTreeMap, BTreeSet, HashMap};
use std::rc::Rc;
Expand Down Expand Up @@ -39,6 +40,7 @@ use crate::query_plan::{
use crate::types::AccuracyTarget;
use planner_types::pre_asap::Source;

mod rate_placement;
mod windows;
pub(super) use windows::gcd;
pub use windows::{prepare_window_implementations, WindowCostModel};
Expand Down Expand Up @@ -96,7 +98,7 @@ impl QueryCompilationInput {
pub(crate) fn retain_physical_candidate(&mut self) -> Result<(), CompileError> {
use asap_physical_operators::physical_planner::{promql_rows, PhysicalCandidate};
let candidate =
promql_rows::compile_fixed_window_rate_aggregation(&self.selected_plan_root)
rate_placement::compile_fixed_window_rate_aggregation(&self.selected_plan_root)
.or_else(|_| {
promql_rows::compile_current_series_readout(&self.selected_plan_root)
.or_else(|_| {
Expand Down Expand Up @@ -377,7 +379,7 @@ impl ScopedAccuracyEvidence {
#[serde(deny_unknown_fields)]
pub struct PhysicalDeploymentContext {
/// Semantic dataset served by this deployment's input channel; never an endpoint.
pub dataset_identity: planner_types::post_asap::LogicalDatasetIdentity,
pub dataset_identity: asap_types::semantic_fragment::LogicalDatasetIdentity,
pub target: PhysicalDeploymentTarget,
#[serde(rename = "collector_ids")]
pub target_collector_ids: Vec<String>,
Expand Down Expand Up @@ -1043,12 +1045,16 @@ impl BackendLocalPlanningInput {
);
proposed
.candidates
.extend(strategy.fixed_window_rate_candidates(&typed).candidates);
proposed.candidates.extend(
strategy
.query_time_rate_aggregation_candidates(&typed)
.candidates,
);
.extend(rate_placement::fixed_window_rate_candidates(
&direct.candidates,
&typed,
));
proposed
.candidates
.extend(rate_placement::query_time_rate_aggregation_candidates(
&direct.candidates,
&typed,
));
proposed.candidates.extend(direct.candidates);
proposed.rejected.extend(direct.rejected);
for candidate in proposed.candidates {
Expand All @@ -1059,7 +1065,7 @@ impl BackendLocalPlanningInput {
.map_err(|error| CompileError::Snapshot(error.to_string()))?;
let compiled = asap_physical_operators::physical_planner::promql_rows::compile_current_series_readout(&root)
.or_else(|_| asap_physical_operators::physical_planner::promql_rows::compile_rate_ranking(&root).map(|(_, program)| program));
if let Ok(physical) = asap_physical_operators::physical_planner::promql_rows::compile_fixed_window_rate_aggregation(&root) {
if let Ok(physical) = rate_placement::compile_fixed_window_rate_aggregation(&root) {
planner_selection_trace.push(serde_json::json!({
"stage":"planner.physical_candidate", "query_id":query.query_id,
"logical_root_id":crate::planner_selection::explained_root_id(&root, &query.accuracy_target),
Expand Down Expand Up @@ -4653,7 +4659,7 @@ fn collect_selected_materializations(
physical_source.as_ref().unwrap_or(node),
None,
composable,
asap_physical_operators::physical_planner::promql_rows::compile_fixed_window_rate_aggregation(node).is_ok(),
rate_placement::compile_fixed_window_rate_aggregation(node).is_ok(),
None,
&mut selected,
)?;
Expand Down Expand Up @@ -5599,7 +5605,7 @@ pub(crate) mod tests {

fn environment(now: u64) -> PhysicalDeploymentContext {
PhysicalDeploymentContext {
dataset_identity: planner_types::post_asap::LogicalDatasetIdentity {
dataset_identity: asap_types::semantic_fragment::LogicalDatasetIdentity {
namespace: "test".into(),
dataset: "metrics".into(),
},
Expand Down
Loading
Loading