From e6e1b420c3842a9cb6991f9148db539da3e5ed7e Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 3 Sep 2026 18:03:00 -0600 Subject: [PATCH] feat(asapquery): compile startup workloads into physical plans --- Cargo.lock | 6 +- control_plane/Cargo.toml | 6 +- control_plane/src/main.rs | 1 + control_plane/src/physical/compiler.rs | 326 +++++++++++++++++- crates/asap_types/Cargo.toml | 2 +- data_plane/Cargo.toml | 6 +- data_plane/src/main.rs | 81 ++++- .../examples/asapquery-planning-snapshot.json | 79 +++++ docs/user_guide/asapquery-profile.md | 27 +- 9 files changed, 498 insertions(+), 36 deletions(-) create mode 100644 docs/examples/asapquery-planning-snapshot.json diff --git a/Cargo.lock b/Cargo.lock index cb41bc184..943595faa 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -343,7 +343,7 @@ dependencies = [ [[package]] name = "asap-aware-mapping" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=739753e33e096c01faccca8e7a1e3da5ad3aab9c#739753e33e096c01faccca8e7a1e3da5ad3aab9c" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=cb50219c582d43f53ab77d3a595bd1ea4a9aa119#cb50219c582d43f53ab77d3a595bd1ea4a9aa119" dependencies = [ "asap-types", "serde", @@ -354,7 +354,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=739753e33e096c01faccca8e7a1e3da5ad3aab9c#739753e33e096c01faccca8e7a1e3da5ad3aab9c" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=cb50219c582d43f53ab77d3a595bd1ea4a9aa119#cb50219c582d43f53ab77d3a595bd1ea4a9aa119" dependencies = [ "asap-types", "promql-parser 0.10.0", @@ -374,7 +374,7 @@ dependencies = [ [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=739753e33e096c01faccca8e7a1e3da5ad3aab9c#739753e33e096c01faccca8e7a1e3da5ad3aab9c" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=cb50219c582d43f53ab77d3a595bd1ea4a9aa119#cb50219c582d43f53ab77d3a595bd1ea4a9aa119" dependencies = [ "serde", "serde_json", diff --git a/control_plane/Cargo.toml b/control_plane/Cargo.toml index bafd44eea..033c0d2a3 100644 --- a/control_plane/Cargo.toml +++ b/control_plane/Cargo.toml @@ -93,8 +93,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 = "739753e33e096c01faccca8e7a1e3da5ad3aab9c" } -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "739753e33e096c01faccca8e7a1e3da5ad3aab9c" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "cb50219c582d43f53ab77d3a595bd1ea4a9aa119" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "cb50219c582d43f53ab77d3a595bd1ea4a9aa119" } # L1 adoption (design-target-architecture.md Part B): the PromQL front # end itself, replacing control_plane's own query_parser/promql.rs. @@ -102,7 +102,7 @@ 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 = "739753e33e096c01faccca8e7a1e3da5ad3aab9c" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "cb50219c582d43f53ab77d3a595bd1ea4a9aa119" } [dev-dependencies] tokio = { version = "1", features = ["full", "test-util"] } diff --git a/control_plane/src/main.rs b/control_plane/src/main.rs index f0f5d804c..8e0ac885b 100644 --- a/control_plane/src/main.rs +++ b/control_plane/src/main.rs @@ -805,6 +805,7 @@ fn compile_physical_plan_request( planner_revision: request.planner_revision, }, physical::compiler::DeploymentEnvironment { + target: physical::compiler::PhysicalDeploymentTarget::DistributedCollectors, collector_ids: request.collector_ids.clone(), capability_snapshot_id: request.capability_snapshot_id, observed_at_unix_ms: now, diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 8983ef4ab..0090f2fab 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -22,8 +22,8 @@ use planner_types::post_asap::{ use planner_types::pre_asap::QueryExpr; use planner_types::workload::{ AccuracyRequirement, DataArrival, DataWorkload, DurationMs, Evidence, EvidenceSource, - Predictability, Query, QueryLanguage, QueryRequirements, QueryTimeScope, QueryWorkload, Rate, - RepeatedDemand, RepeatingEntry, RepetitionInterval, TimeSelection, + Predictability, Query, QueryLanguage, QueryRecurrence, QueryRequirements, QueryTimeScope, + QueryWorkload, Rate, RepeatedDemand, RepeatingEntry, RepetitionInterval, TimeSelection, }; use serde::{Deserialize, Serialize}; use serde_json::{json, Value}; @@ -42,7 +42,7 @@ use crate::query_plan::{ use crate::types_v2::AccuracyTarget; use planner_types::pre_asap::Source; -pub const PLANNER_REVISION: &str = "739753e33e096c01faccca8e7a1e3da5ad3aab9c"; +pub const PLANNER_REVISION: &str = "cb50219c582d43f53ab77d3a595bd1ea4a9aa119"; #[derive(Debug, Clone)] pub struct PlanningQuery { @@ -140,8 +140,10 @@ pub struct TopKMembershipEvidence { pub source: String, } -#[derive(Debug, Clone)] +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(deny_unknown_fields)] pub struct DeploymentEnvironment { + pub target: PhysicalDeploymentTarget, pub collector_ids: Vec, pub capability_snapshot_id: String, pub observed_at_unix_ms: u64, @@ -152,6 +154,39 @@ pub struct DeploymentEnvironment { pub backend_compat: String, } +#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum PhysicalDeploymentTarget { + DistributedCollectors, + BackendLocalRemoteWrite, +} + +/// Versioned startup input for the Collector-free compatibility profile. +/// Query/data semantics use ASAPPlanner's canonical workload types directly; +/// this wrapper adds only backend-owned implementation evidence and lifecycle +/// identity required to choose a concrete physical realization. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +#[serde(deny_unknown_fields)] +pub struct BackendLocalPlanningSnapshot { + pub snapshot_version: u32, + pub query_workload: QueryWorkload, + pub data_workload: DataWorkload, + pub implementation: BackendLocalImplementation, + pub environment: DeploymentEnvironment, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +#[serde(deny_unknown_fields)] +pub struct BackendLocalImplementation { + pub lifecycle_costs: LifecycleCostEvidence, + pub evidence_observed_at_unix_ms: u64, + pub evidence_valid_for_ms: u64, + pub horizon_seconds: f64, + pub window_implementation_id: String, + pub state_layout: String, + pub implementation_cost: ImplementationCostEvidence, +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct PlanEnvelope { pub plan_id: u64, @@ -1418,6 +1453,8 @@ fn authorize_u64_change( #[derive(Debug, Error)] pub enum CompileError { + #[error("invalid backend-local workload snapshot: {0}")] + Snapshot(String), #[error("planner revision mismatch: request={request}, compiler={compiler}")] PlannerRevision { request: String, @@ -1459,6 +1496,138 @@ impl AccuracyEvidenceProvider for QueryEvidence<'_> { #[derive(Debug, Default)] pub struct PhysicalCompiler; +impl BackendLocalPlanningSnapshot { + /// Invoke the pinned Planner from canonical startup workloads and compile + /// one backend-local PhysicalPlan. No CollectorPlan is produced and no + /// precompiled serving artifact is accepted at this boundary. + pub fn compile(self) -> Result { + if self.snapshot_version != 1 { + return Err(CompileError::Snapshot(format!( + "unsupported workload snapshot version {}", + self.snapshot_version + ))); + } + if self.environment.target != PhysicalDeploymentTarget::BackendLocalRemoteWrite { + return Err(CompileError::Snapshot( + "compatibility workload requires backend_local_remote_write target".into(), + )); + } + let mut workload = self.query_workload; + if let Some(embedded) = &workload.data_workload { + if embedded != &self.data_workload { + return Err(CompileError::Snapshot( + "embedded and standalone DataWorkload snapshots disagree".into(), + )); + } + } + workload.data_workload = Some(self.data_workload.clone()); + workload + .validate() + .map_err(|error| CompileError::Snapshot(error.to_string()))?; + if workload.language != QueryLanguage::PromQL { + return Err(CompileError::Snapshot( + "ASAPQuery compatibility profile accepts PromQL workloads only".into(), + )); + } + let ingestion_rate = self + .data_workload + .ingestion_rate + .value_at(self.environment.observed_at_unix_ms) + .copied() + .ok_or_else(|| { + CompileError::Snapshot("DataWorkload requires fresh ingestion_rate evidence".into()) + })?; + let entries = workload.entries().collect::>(); + if entries.is_empty() { + return Err(CompileError::Snapshot( + "QueryWorkload must contain at least one query".into(), + )); + } + let mut queries = Vec::with_capacity(entries.len()); + for (index, entry) in entries.into_iter().enumerate() { + let evaluation_interval_ms = match entry.recurrence { + QueryRecurrence::Repeated(RepeatedDemand::FixedInterval(interval)) => interval.0, + _ => { + return Err(CompileError::Snapshot(format!( + "query {index} must use fixed-interval repeated demand in the MVP profile" + ))) + } + }; + let lookback_ms = entry + .time_selection + .lookback + .ok_or_else(|| { + CompileError::Snapshot(format!("query {index} requires an explicit lookback")) + })? + .0; + if lookback_ms == 0 || lookback_ms % 1_000 != 0 { + return Err(CompileError::Snapshot(format!( + "query {index} lookback must be a positive whole number of seconds" + ))); + } + let accuracy = entry.requirements.accuracy.target(); + let query_string = entry.query.0; + let parsed = + crate::query_parser::parse_query_expr_canonical(&query_string, accuracy.clone()) + .map_err(|error| CompileError::Snapshot(format!("query {index}: {error}")))?; + let metadata = crate::query_parser::qe_to_parsed_query(&parsed); + if metadata.metric_name.is_empty() { + return Err(CompileError::Snapshot(format!( + "query {index} has no unique time-series source" + ))); + } + if !metadata.label_filters.is_empty() { + return Err(CompileError::Snapshot(format!( + "query {index} uses label filters not yet represented by the physical materialization contract" + ))); + } + let lifecycle = LifecyclePlanningInput { + evaluation_interval_ms, + ingestion_rate_per_second: ingestion_rate.0, + evidence_observed_at_unix_ms: self.implementation.evidence_observed_at_unix_ms, + evidence_valid_for_ms: self.implementation.evidence_valid_for_ms, + horizon_seconds: self.implementation.horizon_seconds, + costs: self.implementation.lifecycle_costs.clone(), + }; + let post_asap = select_post_asap(&parsed, accuracy.clone(), &lifecycle, None) + .map_err(|error| CompileError::Snapshot(format!("query {index}: {error}")))?; + let mut cost = self.implementation.implementation_cost.clone(); + cost.workload_fingerprint = + canonical_promql(&query_string).map_err(CompileError::QueryPlan)?; + cost.horizon_seconds = self.implementation.horizon_seconds; + queries.push(PlanningQuery { + query_id: format!("compat-query-{index}"), + query_string, + post_asap, + source: Source::TimeSeries { + metric: metadata.metric_name, + }, + window_secs: lookback_ms / 1_000, + group_by: metadata.group_by_labels, + accuracy, + lifecycle, + window_implementations: vec![WindowImplementationCandidate { + implementation_id: self.implementation.window_implementation_id.clone(), + framework: SummaryWindowFramework::Tumbling, + window_secs: lookback_ms / 1_000, + pane_secs: lookback_ms / 1_000, + state_layout: self.implementation.state_layout.clone(), + cost, + }], + runtime_policy: RuntimeRulePolicy::default(), + }); + } + PhysicalCompiler.compile( + PlanningRequest { + queries, + evidence: HashMap::new(), + planner_revision: PLANNER_REVISION.into(), + }, + self.environment, + ) + } +} + impl PhysicalCompiler { pub fn compile( &self, @@ -1471,6 +1640,14 @@ impl PhysicalCompiler { compiler: PLANNER_REVISION, }); } + if environment.target == PhysicalDeploymentTarget::BackendLocalRemoteWrite + && !environment.collector_ids.is_empty() + { + return Err(CompileError::Query { + query_id: "deployment-target".into(), + reason: "backend-local target cannot declare Collector producers".into(), + }); + } let mut aggregations = Vec::with_capacity(request.queries.len()); let mut readouts = Vec::with_capacity(request.queries.len()); @@ -1547,7 +1724,12 @@ impl PhysicalCompiler { spatial_filter: String::new(), grouping: query.group_by.clone(), item_label: None, - aggregation_input: AggregationInput::SketchEnvelope, + aggregation_input: match environment.target { + PhysicalDeploymentTarget::DistributedCollectors => { + AggregationInput::SketchEnvelope + } + PhysicalDeploymentTarget::BackendLocalRemoteWrite => AggregationInput::Raw, + }, }; let precompute_materialization = backend_plan::aggregation_config_for_materialization(&aggregation)?; @@ -1635,7 +1817,10 @@ impl PhysicalCompiler { output_representation: OutputRepresentation::SummaryState, }); } - let producer_ids = environment.collector_ids.clone(); + let producer_ids = match environment.target { + PhysicalDeploymentTarget::DistributedCollectors => environment.collector_ids.clone(), + PhysicalDeploymentTarget::BackendLocalRemoteWrite => Vec::new(), + }; // Several queries/readouts may intentionally share one maintained // summary. PrecomputePlan is keyed by physical identity, not query ID. let mut materializations_by_fingerprint = BTreeMap::new(); @@ -1647,12 +1832,21 @@ impl PhysicalCompiler { .or_insert(materialization); } let materializations = materializations_by_fingerprint.into_values().collect(); - let precompute_plan = PrecomputePlan::build( - envelope.clone(), - materializations, - &backend_plan, - &producer_ids, - ) + let precompute_plan = match environment.target { + PhysicalDeploymentTarget::DistributedCollectors => PrecomputePlan::build( + envelope.clone(), + materializations, + &backend_plan, + &producer_ids, + ), + PhysicalDeploymentTarget::BackendLocalRemoteWrite => { + PrecomputePlan::build_backend_local( + envelope.clone(), + materializations, + &backend_plan, + ) + } + } .map_err(|error| CompileError::Query { query_id: "precompute-plan".into(), reason: error.to_string(), @@ -2169,6 +2363,7 @@ mod tests { fn environment(now: u64) -> DeploymentEnvironment { DeploymentEnvironment { + target: PhysicalDeploymentTarget::DistributedCollectors, collector_ids: vec!["edge-a".into(), "edge-b".into()], capability_snapshot_id: "caps-7".into(), observed_at_unix_ms: now, @@ -2405,6 +2600,113 @@ mod tests { .expect("valid empty transmission plan"); } + #[test] + fn canonical_workload_snapshot_invokes_backend_local_planning() { + let data_workload = DataWorkload { + arrival: DataArrival::ContinuouslyIngesting, + ingestion_rate: Evidence { + value: Some(Rate(100.0)), + source: EvidenceSource::Declared, + observed_at_ms: None, + valid_for_ms: None, + }, + ..DataWorkload::default() + }; + let query_workload = QueryWorkload { + language: QueryLanguage::PromQL, + query_batch: None, + repeating_queries: Some(vec![RepeatingEntry { + query: Query("quantile_over_time(0.99, m[1m])".into()), + demand: RepeatedDemand::FixedInterval(RepetitionInterval(10_000)), + requirements: QueryRequirements { + accuracy: AccuracyRequirement::Explicit(AccuracyTarget::EpsilonDelta { + epsilon: 0.01, + delta: 0.01, + }), + ..QueryRequirements::default() + }, + predictability: Predictability::Predictable { known_at: None }, + time_selection: TimeSelection { + scope: QueryTimeScope::RealTime, + lookback: Some(DurationMs(60_000)), + as_of: None, + }, + }]), + data_workload: Some(data_workload.clone()), + }; + let mut environment = environment(10_000); + environment.target = PhysicalDeploymentTarget::BackendLocalRemoteWrite; + environment.collector_ids.clear(); + let template = request("template", "quantile_over_time(0.99, m[1m])") + .queries + .remove(0); + let snapshot = BackendLocalPlanningSnapshot { + snapshot_version: 1, + query_workload, + data_workload, + implementation: BackendLocalImplementation { + lifecycle_costs: template.lifecycle.costs, + evidence_observed_at_unix_ms: 9_500, + evidence_valid_for_ms: 60_000, + horizon_seconds: 300.0, + window_implementation_id: "backend-tumbling-v1".into(), + state_layout: "anchored-pane-v1".into(), + implementation_cost: template.window_implementations[0].cost.clone(), + }, + environment, + }; + let first = snapshot + .clone() + .compile() + .expect("first deterministic plan"); + let second = snapshot + .clone() + .compile() + .expect("second deterministic plan"); + assert_eq!(first.envelope, second.envelope); + assert_eq!(first.backend_plan, second.backend_plan); + assert_eq!(first.query_plan, second.query_plan); + assert_eq!(first.transmission_plan, second.transmission_plan); + assert_eq!( + serde_json::to_value(&first.precompute_plan).unwrap(), + serde_json::to_value(&second.precompute_plan).unwrap() + ); + + let encoded = serde_json::to_vec(&snapshot).expect("serialize startup snapshot"); + let decoded: BackendLocalPlanningSnapshot = + serde_json::from_slice(&encoded).expect("deserialize startup snapshot"); + let bundle = decoded.compile().expect("canonical startup planning"); + + assert!(bundle.collector_plans.is_empty()); + assert!(bundle.transmission_plan.rules.is_empty()); + assert_eq!( + bundle.precompute_plan.ingest.protocol, + IngestProtocol::PrometheusRemoteWriteV1 + ); + assert!(bundle.precompute_plan.producers.is_empty()); + assert_eq!(bundle.query_plan.entries.len(), 1); + assert_eq!(bundle.envelope.planner_revision, PLANNER_REVISION); + } + + #[test] + fn checked_in_backend_local_snapshot_is_canonical_and_compilable() { + let source = include_str!("../../../docs/examples/asapquery-planning-snapshot.json"); + let snapshot: BackendLocalPlanningSnapshot = + serde_json::from_str(source).expect("strict canonical workload fixture"); + let encoded = serde_json::to_value(&snapshot).expect("canonical snapshot value"); + let fixture: serde_json::Value = serde_json::from_str(source).expect("fixture JSON"); + assert_eq!(encoded, fixture); + + let plan = snapshot.compile().expect("fixture compiles"); + assert!(plan.collector_plans.is_empty()); + assert!(plan.transmission_plan.rules.is_empty()); + assert_eq!( + plan.precompute_plan.ingest.protocol, + IngestProtocol::PrometheusRemoteWriteV1 + ); + assert_eq!(plan.query_plan.entries.len(), 1); + } + #[test] fn multiple_readouts_share_one_precompute_materialization() { let mut planning_request = request("q-p90", "quantile_over_time(0.90, m[1m])"); diff --git a/crates/asap_types/Cargo.toml b/crates/asap_types/Cargo.toml index 0ea754f36..fe3a0270e 100644 --- a/crates/asap_types/Cargo.toml +++ b/crates/asap_types/Cargo.toml @@ -30,4 +30,4 @@ xxhash-rust = { version = "0.8", features = ["xxh64"] } # 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 = "739753e33e096c01faccca8e7a1e3da5ad3aab9c" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "cb50219c582d43f53ab77d3a595bd1ea4a9aa119" } diff --git a/data_plane/Cargo.toml b/data_plane/Cargo.toml index 89ec7817c..a424258bd 100644 --- a/data_plane/Cargo.toml +++ b/data_plane/Cargo.toml @@ -36,9 +36,9 @@ control_plane = { path = "../control_plane" } # `L4Node` renamed to `SummaryNode` (ASAPPlanner#217); its fields # (`SummaryExpr::Logical(Box)`, `SummaryAgg { col: ColumnRef, # reduction: Reduction, .. }`) are `pre_asap` types, in the same crate now -# (not a separate `asap-ir` import) -- `find_candidates` still needs to -# walk/match them directly. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "739753e33e096c01faccca8e7a1e3da5ad3aab9c" } +# (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 = "cb50219c582d43f53ab77d3a595bd1ea4a9aa119" } # Shared external (workspace) serde.workspace = true diff --git a/data_plane/src/main.rs b/data_plane/src/main.rs index 792c12613..d85a11463 100644 --- a/data_plane/src/main.rs +++ b/data_plane/src/main.rs @@ -58,6 +58,12 @@ struct Args { #[arg(long)] physical_plan: Option, + /// Versioned canonical QueryWorkload + DataWorkload and backend-local + /// implementation evidence. The ASAPQuery profile invokes the pinned + /// Planner and PhysicalCompiler at startup when this is supplied. + #[arg(long)] + planning_snapshot: Option, + /// Cleanup policy for SketchStore retention. /// `circular_buffer`: keep the N most recent windows per agg /// (N comes from each aggregation's `numAggregatesToRetain`). @@ -345,13 +351,22 @@ fn validate_profile(args: &Args) -> Result<()> { if args.streaming_config.is_none() { return Err("the distributed profile requires --streaming-config".into()); } + if args.planning_snapshot.is_some() { + return Err("--planning-snapshot is available only with --profile asapquery".into()); + } return Ok(()); } if args.streaming_config.is_some() { - return Err("--profile asapquery rejects --streaming-config; use --physical-plan".into()); + return Err( + "--profile asapquery rejects --streaming-config; use --planning-snapshot or --physical-plan" + .into(), + ); } - if args.physical_plan.is_none() { - return Err("--profile asapquery requires --physical-plan".into()); + if args.physical_plan.is_some() == args.planning_snapshot.is_some() { + return Err( + "--profile asapquery requires exactly one of --planning-snapshot or --physical-plan" + .into(), + ); } let mut excluded = Vec::new(); if args.enable_otel_ingest { @@ -457,20 +472,45 @@ async fn main() -> Result<()> { ); } - let startup_physical_plan = if let Some(path) = args.physical_plan.as_ref() { + let startup_artifact = if let Some(path) = args.planning_snapshot.as_ref() { let bytes = fs::read(path)?; - let artifact: data_plane::drivers::query::servers::http::PhysicalPlanInstallRequest = + let snapshot: control_plane::physical::compiler::BackendLocalPlanningSnapshot = serde_json::from_slice(&bytes).map_err(|error| { format!( - "failed to decode physical plan artifact {}: {error}", + "failed to decode planning snapshot {}: {error}", path.display() ) })?; + let plan = snapshot + .compile() + .map_err(|error| format!("startup planning failed for {}: {error}", path.display()))?; + Some( + data_plane::drivers::query::servers::http::PhysicalPlanInstallRequest { + precompute_plan: plan.precompute_plan, + transmission_plan: plan.transmission_plan, + backend_plan: plan.backend_plan.encode_to_vec(), + query_plan: plan.query_plan, + storage_routing: None, + adaptation_evidence: Vec::new(), + }, + ) + } else if let Some(path) = args.physical_plan.as_ref() { + let bytes = fs::read(path)?; + Some(serde_json::from_slice(&bytes).map_err(|error| { + format!( + "failed to decode physical plan artifact {}: {error}", + path.display() + ) + })?) + } else { + None + }; + let startup_physical_plan = if let Some(artifact) = startup_artifact { let active = data_plane::drivers::query::servers::http::build_active_physical_plan( artifact, Arc::new(data_plane::storage_engines::types::BackendStorageRouting::empty()), ) - .map_err(|error| format!("invalid physical plan artifact {}: {error}", path.display()))?; + .map_err(|error| format!("invalid startup PhysicalPlan: {error}"))?; if args.profile == RuntimeProfile::Asapquery && (active.backend_plan.plan_id == 0 || !matches!( @@ -1449,6 +1489,33 @@ mod tests { .unwrap_err() .to_string() .contains("rejects --streaming-config")); + + let planned = Args::try_parse_from([ + "data_plane", + "--profile", + "asapquery", + "--planning-snapshot", + "workload.json", + "--forward-unsupported-queries", + ]) + .unwrap(); + assert!(validate_profile(&planned).is_ok()); + + let ambiguous = Args::try_parse_from([ + "data_plane", + "--profile", + "asapquery", + "--planning-snapshot", + "workload.json", + "--physical-plan", + "plan.json", + "--forward-unsupported-queries", + ]) + .unwrap(); + assert!(validate_profile(&ambiguous) + .unwrap_err() + .to_string() + .contains("exactly one")); } // Step-1 of the JSONL deprecation refactor deleted the diff --git a/docs/examples/asapquery-planning-snapshot.json b/docs/examples/asapquery-planning-snapshot.json new file mode 100644 index 000000000..b7a3bcb44 --- /dev/null +++ b/docs/examples/asapquery-planning-snapshot.json @@ -0,0 +1,79 @@ +{ + "snapshot_version": 1, + "query_workload": { + "language": "prom_q_l", + "query_batch": null, + "repeating_queries": [ + { + "query": "quantile_over_time(0.99, m[1m])", + "demand": { "fixed_interval": 10000 }, + "requirements": { + "accuracy": { + "explicit": { + "EpsilonDelta": { "epsilon": 0.01, "delta": 0.01 } + } + }, + "response_latency": "unspecified" + }, + "predictability": { "predictable": { "known_at": null } }, + "time_selection": { + "scope": "real_time", + "lookback": 60000, + "as_of": null + } + } + ], + "data_workload": { + "arrival": "continuously_ingesting", + "ingestion_volume": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, + "ingestion_rate": { "value": 100.0, "source": "declared", "observed_at_ms": null, "valid_for_ms": null }, + "input_cardinality": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, + "distribution": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null } + } + }, + "data_workload": { + "arrival": "continuously_ingesting", + "ingestion_volume": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, + "ingestion_rate": { "value": 100.0, "source": "declared", "observed_at_ms": null, "valid_for_ms": null }, + "input_cardinality": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, + "distribution": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null } + }, + "implementation": { + "lifecycle_costs": { + "build": 10.0, + "maintenance_per_update": 0.001, + "read": 0.1, + "retention_per_second": 0.001, + "retirement": 1.0 + }, + "evidence_observed_at_unix_ms": 9500, + "evidence_valid_for_ms": 60000, + "horizon_seconds": 300.0, + "window_implementation_id": "backend-tumbling-v1", + "state_layout": "anchored-pane-v1", + "implementation_cost": { + "model_version": "compat-cost-v1", + "workload_fingerprint": "replaced-with-canonical-promql", + "observed_at_unix_ms": 9500, + "valid_for_ms": 60000, + "horizon_seconds": 300.0, + "cpu_cost": 1.0, + "peak_memory_bytes": 1024, + "network_bytes": 0, + "storage_bytes": 512, + "source_scan_bytes": 0, + "weighted_cost": 1.0 + } + }, + "environment": { + "target": "backend_local_remote_write", + "collector_ids": [], + "capability_snapshot_id": "asapquery-local-v1", + "observed_at_unix_ms": 10000, + "max_evidence_age_ms": 60000, + "plan_version": 1, + "activation_unix_ms": 10000, + "expiry_unix_ms": null, + "backend_compat": "asap-query-backend.v1" + } +} diff --git a/docs/user_guide/asapquery-profile.md b/docs/user_guide/asapquery-profile.md index 1ded9be78..7e637e4e0 100644 --- a/docs/user_guide/asapquery-profile.md +++ b/docs/user_guide/asapquery-profile.md @@ -17,7 +17,7 @@ Then start the backend: ```bash cargo run -p data_plane -- \ --profile asapquery \ - --physical-plan /etc/asapquery/physical-plan.json \ + --planning-snapshot /etc/asapquery/workload.json \ --prometheus-server http://127.0.0.1:9090 \ --forward-unsupported-queries \ --http-port 9091 \ @@ -54,7 +54,21 @@ Receiver evidence is exported from `/metrics` as `asap_remote_write_rejected_requests_total`, and `asap_remote_write_bytes_total`. -`physical-plan.json` is the JSON representation accepted by +`workload.json` is a versioned `BackendLocalPlanningSnapshot`. Its +`query_workload` and `data_workload` fields deserialize directly into the +canonical ASAPPlanner types; `implementation` contains only backend-owned +cost/window evidence, and `environment.target` must be +`backend_local_remote_write`. Startup validates the workload, calls the pinned +ASAPPlanner, compiles matching PrecomputePlan/BackendPlan/QueryPlan views, and +installs the resulting immutable snapshot before accepting traffic. It does +not create or wait for a CollectorPlan. + +A complete canonical input is checked in at +`docs/examples/asapquery-planning-snapshot.json`. + +For reproducibility/debugging, `--physical-plan` remains an alternative to +`--planning-snapshot`; exactly one is required. `physical-plan.json` is the +JSON representation accepted by `POST /api/v1/physical-plan`: a `PrecomputePlan`, `TransmissionPlan`, encoded `BackendPlan`, and authoritative `QueryPlan` DAG with one shared plan identity and version. For this profile, the precompute ingest contract must be @@ -70,8 +84,7 @@ Query serving uses only the installed `QueryPlan` DAG. A query absent from that DAG is a capability miss and goes to the exact Prometheus fallback; the backend does not search materialization candidates while serving. -Automatic startup compilation from canonical `QueryWorkload` and -`DataWorkload` snapshots is not part of this phase. Until the workload types -have a stable serialized contract and the physical compiler supports a -backend-only deployment target, the artifact must be produced offline by the -planner/control-plane pipeline. +Snapshot schema version `1` currently accepts fixed-interval repeating PromQL +queries with explicit whole-second lookbacks and fresh ingestion-rate +evidence. Unsupported snapshot semantics fail startup rather than silently +inventing cost or placement evidence.