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
7 changes: 4 additions & 3 deletions Cargo.lock

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

6 changes: 3 additions & 3 deletions control_plane/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -93,16 +93,16 @@ 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 = "https://github.com/ProjectASAP/ASAPPlanner", rev = "3afcba68f4e8397fb81e2be988f47120f63f7a39" }
asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "3afcba68f4e8397fb81e2be988f47120f63f7a39" }
planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "5d0b6f6edcac65edc89a72051f37977ab0c83031" }
asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "5d0b6f6edcac65edc89a72051f37977ab0c83031" }

# L1 adoption (design-target-architecture.md Part B): the PromQL front
# end itself, replacing control_plane's own query_parser/promql.rs.
# Pinned via `rev`, not a floating branch reference. Same rev as
# `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 = "https://github.com/ProjectASAP/ASAPPlanner", rev = "3afcba68f4e8397fb81e2be988f47120f63f7a39" }
asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "5d0b6f6edcac65edc89a72051f37977ab0c83031" }

[dev-dependencies]
tokio = { version = "1", features = ["full", "test-util"] }
Expand Down
8 changes: 8 additions & 0 deletions control_plane/proto/backend_plan.proto
Original file line number Diff line number Diff line change
Expand Up @@ -167,6 +167,14 @@ message Materialization {
SummaryParams params = 6;
ColumnRef col = 7;
optional RetentionPolicy retention = 8;
optional SummaryMaintenanceLifecycle lifecycle = 9;
}

message SummaryMaintenanceLifecycle {
string kind = 1;
string maintenance_mode = 2;
string evaluation_schedule = 3;
string output_representation = 4;
}

message RoutingEntry {
Expand Down
6 changes: 3 additions & 3 deletions control_plane/src/asap_tier_analysis.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,10 +49,10 @@ use std::time::Duration;

use promql_parser::parser::{self, Expr, VectorSelector};

use crate::intent_algebra::agg_intent::AggIntent;
use crate::intent_algebra::query_expr::QueryExpr;
use crate::query_parser::{parse_query_expr_canonical, parsed_query_from_canonical};
use crate::types_v2::AccuracyTarget;
use planner_types::pre_asap::AggIntent;
use planner_types::pre_asap::QueryExpr;

pub use crate::sketch_algebra::capability::{
capability_for, Capability, OuterAgg, OuterFn, SketchKindHandle,
Expand Down Expand Up @@ -392,7 +392,7 @@ pub(crate) fn collect_agg_intents(expr: &QueryExpr, out: &mut Vec<AggIntent>) {
/// the AST walker can recover them; this fallback runs when the AST
/// walk fails.
fn intent_kind_label(intent: &AggIntent) -> &'static str {
if crate::intent_algebra::as_frequency(intent).is_some() {
if crate::planner_selection::as_frequency(intent).is_some() {
return "frequency";
}
match intent {
Expand Down
2 changes: 1 addition & 1 deletion control_plane/src/asap_tier_implement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -76,9 +76,9 @@ use std::rc::Rc;
use asap_aware_mapping::DefaultCostModel;
use planner_types::post_asap::SummaryNode;

use crate::intent_algebra::query_expr::QueryExpr;
use crate::query_parser::parse_query_expr_canonical;
use crate::types_v2::AccuracyTarget;
use planner_types::pre_asap::QueryExpr;

/// Fixed accuracy target for this L1 call site (L1 adoption,
/// design-target-architecture.md Part B) -- matches
Expand Down
4 changes: 3 additions & 1 deletion control_plane/src/backend_plan/from_stage_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,9 +25,10 @@ use asap_types::{AggregationConfig, MonitorSpec, PolicyFingerprint, QueryLanguag

use crate::emit::monitor::{agg_id_for_metric, MonitorIntent};
use crate::emit::stage_config::build_backend_aggregation_json;
use crate::intent_algebra::{ColumnRef, Source, WindowKind};
use crate::physical::colored_dag::emitter::{BackendAggregation, BackendStageConfig};
use crate::sketch_algebra::capability::{Capability, SketchKindHandle};
use asap_types::enums::WindowKind;
use planner_types::pre_asap::{ColumnRef, Source};

use super::{BackendPlan, Materialization, RoutingEntry, StorageBackend, WindowSpec};

Expand Down Expand Up @@ -74,6 +75,7 @@ pub fn from_stage_config(
params,
col: ColumnRef::SampleValue,
retention: None,
lifecycle: None,
},
);
}
Expand Down
77 changes: 76 additions & 1 deletion control_plane/src/backend_plan/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -38,8 +38,9 @@ use prost::Message as _;
use serde::{Deserialize, Serialize};
use thiserror::Error;

use crate::intent_algebra::{ColumnRef, Source, WindowKind};
use crate::sketch_algebra::capability::{Capability, SketchKindHandle};
use asap_types::enums::WindowKind;
use planner_types::pre_asap::{ColumnRef, Source};

/// Errors decoding a `BackendPlan` (or one of its parts) from its proto
/// wire form. Encoding (`From<&T> for proto::T`) is always infallible —
Expand All @@ -55,6 +56,8 @@ pub enum DecodeError {
MissingOneof(&'static str),
#[error("unspecified/unknown enum value {value} for {field}")]
UnknownEnumValue { field: &'static str, value: i32 },
#[error("unsupported summary-maintenance lifecycle: {0}")]
UnsupportedLifecycle(String),
}

// ── WindowSpec ───────────────────────────────────────────────────────────────
Expand Down Expand Up @@ -575,6 +578,37 @@ pub struct Materialization {
pub params: SummaryParams,
pub col: ColumnRef,
pub retention: Option<RetentionPolicy>,
pub lifecycle: Option<SummaryMaintenanceLifecycle>,
}

#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SummaryMaintenanceLifecycle {
pub kind: String,
pub maintenance_mode: String,
pub evaluation_schedule: String,
pub output_representation: String,
}

impl From<&SummaryMaintenanceLifecycle> for proto::SummaryMaintenanceLifecycle {
fn from(value: &SummaryMaintenanceLifecycle) -> Self {
Self {
kind: value.kind.clone(),
maintenance_mode: value.maintenance_mode.clone(),
evaluation_schedule: value.evaluation_schedule.clone(),
output_representation: value.output_representation.clone(),
}
}
}

impl From<proto::SummaryMaintenanceLifecycle> for SummaryMaintenanceLifecycle {
fn from(value: proto::SummaryMaintenanceLifecycle) -> Self {
Self {
kind: value.kind,
maintenance_mode: value.maintenance_mode,
evaluation_schedule: value.evaluation_schedule,
output_representation: value.output_representation,
}
}
}

impl From<&Materialization> for proto::Materialization {
Expand All @@ -588,6 +622,7 @@ impl From<&Materialization> for proto::Materialization {
params: Some((&m.params).into()),
col: Some((&m.col).into()),
retention: m.retention.as_ref().map(Into::into),
lifecycle: m.lifecycle.as_ref().map(Into::into),
}
}
}
Expand All @@ -599,6 +634,22 @@ impl TryFrom<proto::Materialization> for Materialization {
m.params
.ok_or(DecodeError::MissingOneof("Materialization.params"))?,
)?;
let lifecycle = m.lifecycle.map(SummaryMaintenanceLifecycle::from);
if let Some(lifecycle) = &lifecycle {
if lifecycle.kind != "continuously_maintained"
|| lifecycle.maintenance_mode != "incremental"
|| lifecycle.evaluation_schedule != "per_update"
|| lifecycle.output_representation != "summary_state"
{
return Err(DecodeError::UnsupportedLifecycle(format!(
"{}/{}/{}/{}",
lifecycle.kind,
lifecycle.maintenance_mode,
lifecycle.evaluation_schedule,
lifecycle.output_representation
)));
}
}
Ok(Materialization {
fingerprint: PolicyFingerprint(m.fingerprint),
source: m
Expand All @@ -618,6 +669,7 @@ impl TryFrom<proto::Materialization> for Materialization {
.ok_or(DecodeError::MissingOneof("Materialization.col"))?
.try_into()?,
retention: m.retention.map(Into::into),
lifecycle,
})
}
}
Expand Down Expand Up @@ -760,6 +812,12 @@ mod tests {
retention: Some(RetentionPolicy {
num_aggregates_to_retain: Some(1000),
}),
lifecycle: Some(SummaryMaintenanceLifecycle {
kind: "continuously_maintained".into(),
maintenance_mode: "incremental".into(),
evaluation_schedule: "per_update".into(),
output_representation: "summary_state".into(),
}),
}
}

Expand Down Expand Up @@ -807,6 +865,23 @@ mod tests {
}
}

#[test]
fn unsupported_lifecycle_fails_closed_on_decode() {
let mut plan = sample_plan();
plan.materializations
.values_mut()
.next()
.unwrap()
.lifecycle
.as_mut()
.unwrap()
.kind = "ephemeral".into();
assert!(matches!(
BackendPlan::decode(&plan.encode_to_vec()),
Err(DecodeError::UnsupportedLifecycle(_))
));
}

#[test]
fn round_trips_a_plan_with_both_exact_and_approximate_materializations() {
let plan = sample_plan();
Expand Down
2 changes: 1 addition & 1 deletion control_plane/src/emit/stage_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -62,9 +62,9 @@ use crate::physical::colored_dag::emitter::{
// archive-tier metric lists). Importing them at module scope produced an
// unused-import warning on every non-test build, so they're scoped into the
// test module's `use super::*` instead (P2-5).
use crate::intent_algebra::ColumnRef;
use crate::physical::colored_dag::stage_id::StageId;
use planner_types::post_asap::{SketchAlgorithm, SketchParams, SketchQuery};
use planner_types::pre_asap::ColumnRef;
// `BackendAggregation.sketch_kind`/`.sketch_params` span both exact
// accumulators and approximate sketches -- see
// `physical::colored_dag::emitter`'s `use asap_types::{...}` note.
Expand Down
Loading