From 6b79822e74dd6ffb35d60b977cd2ae436c8eb0c7 Mon Sep 17 00:00:00 2001 From: zzylol Date: Tue, 29 Sep 2026 01:25:44 +0000 Subject: [PATCH] feat(control-plane): compute ERP and analytical workload costs Without provider quotes, compute complete workload costs for each physically compiled candidate from workload demand, applicable ERP measurements and analytical resource estimates, including native operator CPU, heap workspace and maintenance programs. Report each candidate's resource breakdown in cost_comparison; complete provider quotes remain an optional override. Reject candidates that fail ERP accuracy admission before pricing, and accept a bounded classic-HLL confidence contract in scoped evidence for backend-local Regular HLL readouts. Split from #761 so #728's structural checks precede costing. The tree equals the pre-split #728 head. Co-Authored-By: Claude Opus 5.5 --- README.md | 17 +- control_plane/src/main.rs | 11 +- control_plane/src/physical/compiler.rs | 218 ++++- control_plane/src/physical/erp.rs | 175 +++- .../src/physical/post_asap/cost_model.rs | 251 +++++- control_plane/src/physical/workload_cost.rs | 282 ++++-- .../src/physical/workload_cost/automatic.rs | 810 ++++++++++++++++++ control_plane/tests/native_rate_topk.rs | 67 +- control_plane/tests/native_snapshot_topk.rs | 45 + .../support/distinct_planning_process.rs | 179 ++++ docs/design_docs/asapplanner-integration.md | 7 + .../evidence-dependent-candidates.md | 208 ++++- docs/design_docs/shape-aware-erp-v1.md | 8 +- .../catalog-physical-plan-runtime.md | 3 +- docs/evaluation/e2e-physical-dag.md | 18 +- docs/examples/workload-cost-evidence.md | 18 +- 16 files changed, 2191 insertions(+), 126 deletions(-) create mode 100644 control_plane/src/physical/workload_cost/automatic.rs diff --git a/README.md b/README.md index f9391cbc1..ae06b77f7 100644 --- a/README.md +++ b/README.md @@ -302,13 +302,16 @@ dot -Tsvg target/readme-evidence/selected.dot \ -o target/readme-evidence/selected.svg ``` -Candidate discovery accepts the checked-in unquoted templates. Deployment and -selected-plan inspection require `ASAPQUERY_PLANNING_SNAPSHOT` to point to a -snapshot with complete, valid workload cost evidence. Prepare that input using -the [cost evidence workflow](docs/examples/workload-cost-evidence.md). -There is one snapshot compiler: it compares complete executable alternatives, -including exact fallback. Materialization IDs are definitions, not physical SIDs. -For `--metricsql`, collect quotes for the MetricsQL frontend. +Candidate discovery and deployment accept snapshots without external workload +quotes. The backend combines applicable ERP resources or analytical estimates +with data size, query frequency and physical-plan structure, then compares +complete candidate costs, including exact fallback. The selected plan includes +the resource breakdown and assumptions in `cost_comparison`. + +Point `ASAPQUERY_PLANNING_SNAPSHOT` to your workload snapshot with current data +and capability inputs. Optional calibrated provider quotes can override automatic +costing through the [cost evidence workflow](docs/examples/workload-cost-evidence.md). +Materialization IDs are definitions, not physical SIDs. ## Prometheus runbook diff --git a/control_plane/src/main.rs b/control_plane/src/main.rs index 74a185f00..8464ec173 100644 --- a/control_plane/src/main.rs +++ b/control_plane/src/main.rs @@ -731,7 +731,7 @@ fn compile_physical_plan_request( ) } }, - None => frontend.compile(compilation_request, environment), + None => physical::workload_cost::select_candidates(candidates, environment, None, frontend), }; let bundle = match compiled { Ok(bundle) => bundle, @@ -986,6 +986,15 @@ mod api_tests { let (plan, collectors, _, _, _) = compile_physical_plan_request(request, false, QueryFrontend::PromQl).unwrap(); let plan = plan.unwrap(); + let report = plan + .cost_comparison + .as_ref() + .expect("HTTP must compare automatically priced candidates"); + assert_eq!(report.model_version, "backend-workload-resources-v2"); + assert!(report + .candidate_evaluations + .iter() + .any(|candidate| candidate.automatic_cost.is_some())); assert!(collectors.is_empty()); assert!(plan.collector_plans.is_empty()); assert_eq!( diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 11f892647..1abb20168 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -228,6 +228,8 @@ pub struct ScopedAccuracyEvidence { #[serde(default)] pub quantile_operand_domains: Vec, #[serde(default)] + pub hll: Option, + #[serde(default)] pub values_non_negative: Option, #[serde(default)] pub input_row_count: Option, @@ -284,6 +286,10 @@ impl ScopedAccuracyEvidence { && self .topk_max_distinct_items .is_none_or(|n| n > 0 && n <= (1_u64 << 53)) + && self + .hll + .as_ref() + .is_none_or(super::erp::HllConfidenceContract::valid) && self.input_row_count != Some(0) && self .hydra_shared_grid_collision_bound @@ -364,6 +370,7 @@ pub enum PhysicalDeploymentTarget { pub struct BackendLocalPlanningInput { #[serde(rename = "snapshot_version")] pub schema_version: u32, + /// Optional complete provider override; absent evidence uses backend ERP/analytical workload costing. #[serde(default, skip_serializing_if = "Option::is_none")] pub workload_cost_evidence: Option, #[serde(deserialize_with = "deserialize_snapshot_query_workload")] @@ -745,23 +752,16 @@ impl BackendLocalPlanningInput { self, frontend: QueryFrontend, ) -> Result { - let evidence = self.workload_cost_evidence.clone().ok_or_else(|| { - CompileError::Snapshot( - "deployment requires complete workload cost evidence; export candidates and price them before compiling".into(), - ) - })?; + let evidence = self.workload_cost_evidence.clone(); let (request, environment) = self.into_physical_compilation_request()?; let candidates = super::workload_cost::enumerate_exact_and_materialized_candidates(request)?; - if frontend == QueryFrontend::MetricsQl { - super::workload_cost::select_lowest_cost_metricsql_candidate( - candidates, - environment, - &evidence, - ) - } else { - super::workload_cost::select_lowest_cost_candidate(candidates, environment, &evidence) - } + super::workload_cost::select_candidates( + candidates, + environment, + evidence.as_ref(), + frontend, + ) } /// Build Planner-authorized candidates for evidence collection without publishing. @@ -1441,7 +1441,12 @@ impl DeploymentPlanCompiler { validate_evidence(&query.query_id, e, &environment)?; } let node = query.selected_plan_root.clone(); - reject_uncertified_readouts(&query.query_id, &node, &query.accuracy_target)?; + reject_uncertified_readouts( + &query.query_id, + &node, + &query.accuracy_target, + environment.target, + )?; let selected = if super::maintained_population::supported_node(&node) { Vec::new() } else { @@ -1455,7 +1460,12 @@ impl DeploymentPlanCompiler { })? }; if !selected.is_empty() { - reject_uncertified_readouts(&query.query_id, &node, &query.accuracy_target)?; + reject_uncertified_readouts( + &query.query_id, + &node, + &query.accuracy_target, + environment.target, + )?; } let selected = selected .into_iter() @@ -2794,13 +2804,8 @@ fn logical_roots_and_candidates( } } for (accuracy, scope, roots) in cohorts { - // ERP v1 has no calibrated failure probability. Preserve explicit - // confidence requirements through theoretical/exact fallback. let scoped_erp = erp.map(|policy| { let mut policy = policy.clone(); - if !matches!(accuracy, AccuracyTarget::Epsilon(_)) { - policy.artifact.records.clear(); - } if policy.observed_populations.is_some() && !roots .iter() @@ -2823,7 +2828,15 @@ fn logical_roots_and_candidates( }); policy }); - let erp = scoped_erp.as_ref(); + // Confidence limitations invalidate an empirical accuracy decision, + // not the independently matched resource measurements. + let accuracy_erp = scoped_erp.clone().map(|mut policy| { + if !matches!(accuracy, AccuracyTarget::Epsilon(_)) { + policy.artifact.records.clear(); + } + policy + }); + let erp = accuracy_erp.as_ref(); let mut model = ControlPlaneCostModel::new(accuracy.clone()).with_exact_composition_costs( scope .as_ref() @@ -2834,9 +2847,17 @@ fn logical_roots_and_candidates( if let Some(erp) = erp { model = model.with_erp(erp.clone()); } + if let Some(costs) = &scoped_erp { + model = model.with_erp_costs(costs.clone()); + } let certificate = scope.as_ref().and_then(|id| evidence.get(id)); let scoped_certificate = scope.as_ref().and_then(|id| scoped_evidence.get(id)); + let hll = scoped_certificate.and_then(|evidence| evidence.hll.as_ref()); + if let Some(contract) = hll { + model = model.with_hll_confidence(contract.clone()); + } let accuracy_model = super::erp::ErpAccuracyModel { + hll, policy: erp, max_error: match accuracy { AccuracyTarget::Epsilon(e) | AccuracyTarget::EpsilonDelta { epsilon: e, .. } => e, @@ -2891,6 +2912,7 @@ fn logical_roots_and_candidates( "source": evidence.source, "data_snapshot_id": evidence.data_snapshot_id, "observed_at_unix_ms": evidence.observed_at_unix_ms, + "hll": evidence.hll, }); } trace["deployment_overrides"] = serde_json::json!([]); @@ -2921,7 +2943,7 @@ fn logical_roots_and_candidates( Ok(traces) } -fn requires_exact_erp_fallback( +pub(super) fn requires_exact_erp_fallback( node: &SummaryNode, accuracy: &AccuracyTarget, erp: &super::erp::ErpPlanningInput, @@ -3984,6 +4006,7 @@ fn reject_uncertified_readouts( query_id: &str, root: &Rc, accuracy: &AccuracyTarget, + target: PhysicalDeploymentTarget, ) -> Result<(), CompileError> { let dag = planner_types::post_asap::compile_executable_dag(root).map_err(|error| { CompileError::Query { @@ -3992,6 +4015,21 @@ fn reject_uncertified_readouts( } })?; for node in &dag.nodes { + if target != PhysicalDeploymentTarget::BackendLocalRemoteWrite + && node.guarantee.as_ref().is_some_and(|guarantee| { + guarantee.provenance.iter().any(|source| { + matches!(source, + planner_types::post_asap::GuaranteeSource::SketchReadout { contract, .. } + if contract == "classic_hll_linear_counting_collision_bound_v1") + }) + }) + { + return Err(CompileError::Query { + query_id: query_id.into(), + reason: "classic HLL confidence is bound to the backend-local Regular estimator" + .into(), + }); + } if matches!( node.payload, planner_types::post_asap::ExecutableOperatorPayload::SummaryEstimate { .. } @@ -5789,6 +5827,97 @@ pub(crate) mod tests { assert!(result.unwrap().precompute_plan.materializations.is_empty()); } + fn bounded_hll_snapshot_wire() -> serde_json::Value { + let mut wire = serde_json::to_value(planning_snapshot()).unwrap(); + let query = "distinct_over_time(data[5s])"; + wire["query_workload"]["repeating_queries"][0]["query"] = query.into(); + wire["query_workload"]["repeating_queries"][0]["requirements"]["accuracy"] = + serde_json::json!({"explicit":{"EpsilonDelta":{"epsilon":0.05,"delta":0.01}}}); + wire["implementation"]["data_snapshot_id"] = "hll-bounded-population".into(); + let now = wire["environment"]["observed_at_unix_ms"].clone(); + wire["implementation"]["accuracy_evidence"] = serde_json::json!({query: { + "query_string": query, "data_snapshot_id": "hll-bounded-population", + "data_workload": wire["data_workload"], "source": "enforced-distinct-domain-v1", + "observed_at_unix_ms": now, "valid_for_ms": 60000, + "hll": {"model":"asap-classic64-uniform-hash-linear-counting-v1", + "max_distinct_per_readout": 128} + }}); + wire + } + + /// An applicable classic-HLL confidence contract restores normal selection. + #[test] + fn bounded_hll_confidence_selects_and_binds_a_materialization() { + let wire = bounded_hll_snapshot_wire(); + let snapshot: BackendLocalPlanningInput = serde_json::from_value(wire).unwrap(); + let (request, environment) = snapshot.into_physical_compilation_request().unwrap(); + let plan = DeploymentPlanCompiler + .compile_promql(request, environment) + .unwrap(); + assert_eq!(plan.precompute_plan.materializations.len(), 1, "{plan:#?}"); + assert_eq!( + plan.precompute_plan.materializations[0].aggregation_type, + asap_types::AggregationType::HLL + ); + } + + /// A backend-local estimator proof cannot certify an unverified collector implementation. + #[test] + fn hll_confidence_cannot_be_rebound_to_collectors() { + let snapshot: BackendLocalPlanningInput = + serde_json::from_value(bounded_hll_snapshot_wire()).unwrap(); + let (request, _) = snapshot.into_physical_compilation_request().unwrap(); + let error = reject_uncertified_readouts( + "q", + &request.queries[0].selected_plan_root, + &request.queries[0].accuracy_target, + PhysicalDeploymentTarget::DistributedCollectors, + ) + .unwrap_err(); + assert!(error + .to_string() + .contains("backend-local Regular estimator")); + } + + /// Missing contracts and unattainable targets cannot acquire a confidence proof. + #[test] + fn hll_confidence_keeps_exact_when_absent_or_insufficient() { + for missing in [true, false] { + let mut wire = bounded_hll_snapshot_wire(); + if missing { + wire["implementation"]["accuracy_evidence"] = serde_json::json!({}); + } else { + wire["query_workload"]["repeating_queries"][0]["requirements"]["accuracy"] = + serde_json::json!({"explicit":{"EpsilonDelta":{"epsilon":0.05,"delta":1e-12}}}); + } + let snapshot: BackendLocalPlanningInput = serde_json::from_value(wire).unwrap(); + let (request, environment) = snapshot.into_physical_compilation_request().unwrap(); + let plan = DeploymentPlanCompiler + .compile_promql(request, environment) + .unwrap(); + assert!( + plan.precompute_plan.materializations.is_empty(), + "{plan:#?}" + ); + } + } + + /// An estimator mismatch or unsupported population must be rejected at admission. + #[test] + fn hll_confidence_rejects_invalid_source_contracts() { + for contract in [ + serde_json::json!({"model":"asap-classic64-uniform-hash-linear-counting-v1","max_distinct_per_readout":0}), + serde_json::json!({"model":"asap-classic64-uniform-hash-linear-counting-v1","max_distinct_per_readout":4097}), + serde_json::json!({"model":"hip-rse","max_distinct_per_readout":128}), + ] { + let mut wire = bounded_hll_snapshot_wire(); + wire["implementation"]["accuracy_evidence"]["distinct_over_time(data[5s])"]["hll"] = + contract; + let snapshot: BackendLocalPlanningInput = serde_json::from_value(wire).unwrap(); + assert!(snapshot.into_physical_compilation_request().is_err()); + } + } + /// Accuracy facts must belong to the same query, workload, snapshot and /// evidence window before they enter Planner. #[test] @@ -5804,6 +5933,7 @@ pub(crate) mod tests { source: "enforced-source-contract".into(), observed_at_unix_ms: 9_000, valid_for_ms: 2_000, + hll: None, quantile_operand_domains: vec![], values_non_negative: Some(true), input_row_count: Some(10), @@ -5888,6 +6018,7 @@ pub(crate) mod tests { source: "enforced-source-contract".into(), observed_at_unix_ms: 9_000, valid_for_ms: 2_000, + hll: None, quantile_operand_domains: vec![QuantileOperandDomainEvidence { operand: serde_json::to_value(lhs.as_ref()).unwrap(), lower: 1.0, @@ -7116,6 +7247,29 @@ pub(crate) mod tests { .unwrap() } + /// ERP memory remains usable under an explicit confidence target, even + /// though ERP v1 error observations cannot certify that target. + #[test] + fn erp_resource_costs_survive_confidence_target() { + let mut snapshot = planning_snapshot(); + let entry = &mut snapshot.query_workload.repeating_queries.as_mut().unwrap()[0]; + entry.query = Query("quantile_over_time(0.9,m[1m])".into()); + entry.requirements.accuracy = AccuracyRequirement::Explicit(AccuracyTarget::EpsilonDelta { + epsilon: 0.1, + delta: 0.01, + }); + snapshot.physical_inputs.erp = + Some(crate::physical::post_asap::cost_model::tests::erp_cost_fixture()); + let (request, _) = snapshot.into_physical_compilation_request().unwrap(); + assert!(request + .planner_selection_trace + .iter() + .flat_map(|trace| trace["groups"].as_array().unwrap()) + .flat_map(|group| group["candidates"].as_array().unwrap()) + .any(|candidate| candidate["cost_estimate"]["source"] == "erp" + && candidate["estimated_cost"] == 7777.0)); + } + fn derived_query_window_secs(query: &str) -> u64 { let mut snapshot = planning_snapshot(); let entry = &mut snapshot.query_workload.repeating_queries.as_mut().unwrap()[0]; @@ -8448,8 +8602,13 @@ pub(crate) mod tests { assert_eq!(encoded, fixture); assert!( - snapshot.clone().compile_promql().is_err(), - "discovery fixtures must be priced before deployment" + snapshot + .clone() + .compile_promql() + .unwrap() + .cost_comparison + .is_some(), + "deployment computes complete workload costs automatically" ); let (local, env) = snapshot .clone() @@ -8488,8 +8647,13 @@ pub(crate) mod tests { let snapshot: BackendLocalPlanningInput = serde_json::from_str(source).expect("strict compatibility demo fixture"); assert!( - snapshot.clone().compile_promql().is_err(), - "discovery fixtures must be priced before deployment" + snapshot + .clone() + .compile_promql() + .unwrap() + .cost_comparison + .is_some(), + "deployment computes complete workload costs automatically" ); let (local, env) = snapshot .clone() diff --git a/control_plane/src/physical/erp.rs b/control_plane/src/physical/erp.rs index 72dfb3e2a..1cbcb6681 100644 --- a/control_plane/src/physical/erp.rs +++ b/control_plane/src/physical/erp.rs @@ -204,10 +204,44 @@ impl ReadoutEvidence { } } -/// ERP v1 measures error magnitudes, not tail probabilities. Only an explicit -/// epsilon-only request may use these observations as its accuracy contract. +/// Trusted source contract for every HLL readout population of one scoped query. +/// The upper bound covers the union of all merged panes, not each pane separately. +/// This is not derived from observed series count or a sampled distinct count. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +#[serde(deny_unknown_fields)] +pub struct HllConfidenceContract { + /// Identifies classic 64-bit HLL with its linear-counting correction, under + /// independent uniform bucket hashing. Other estimators are not certified. + pub model: String, + pub max_distinct_per_readout: u32, +} + +impl HllConfidenceContract { + pub(crate) fn valid(&self) -> bool { + self.model == "asap-classic64-uniform-hash-linear-counting-v1" + && (1..=4096).contains(&self.max_distinct_per_readout) + } + + pub(crate) fn confidence( + &self, + epsilon: f64, + ) -> Option { + self.valid() + .then(|| { + asap_aware_mapping::hll_confidence::ClassicHllConfidence::new( + self.max_distinct_per_readout, + epsilon, + ) + }) + .flatten() + } +} + +/// ERP error magnitudes and estimator-specific confidence are separate inputs. +/// ERP v1 observed maxima do not establish a failure-probability guarantee. pub(crate) struct ErpAccuracyModel<'a> { pub policy: Option<&'a ErpPlanningInput>, + pub hll: Option<&'a HllConfidenceContract>, pub max_error: f64, } @@ -224,6 +258,20 @@ impl asap_aware_mapping::AccuracyModel for ErpAccuracyModel<'_> { query: &planner_types::post_asap::SketchQuery, ) -> Option { use planner_types::post_asap::*; + if let ( + Some(contract), + SummaryFamilyType::Sketch(kind, GroupingStrategy::PerSubpopulationInstance), + SketchQuery::Cardinality, + ) = (self.hll, family, query) + { + if let (SketchAlgorithm::Hll, SketchParams::Hll { precision }) = + (kind.algorithm(), kind.params()) + { + // Resource ERP can still rank this state. Error maxima must not + // replace the independent, estimator-specific confidence bound. + return contract.confidence(self.max_error)?.guarantee(*precision); + } + } let theoretical = asap_aware_mapping::DefaultAccuracyModel.local_guarantee(family, query); let (Some(policy), SummaryFamilyType::Sketch(kind, _)) = (self.policy, family) else { return theoretical; @@ -575,6 +623,91 @@ impl ErpParameterDecision { } impl ErpPlanningInput { + /// Match ERP resource measurements for an already-sized candidate. Error + /// magnitudes only identify existing records here: no accuracy target is + /// certified by this cost lookup. Reuse ERP's implementation, population, + /// shape, minimum-trial and runtime checks rather than a second matcher. + pub(crate) fn candidate_resources( + &self, + algorithm: &SketchAlgorithm, + params: &SketchParams, + ) -> Option<(asap_aware_mapping::erp::ErpResourceProfile, Vec)> { + self.artifact.validate().ok()?; + let metrics = self + .artifact + .records + .iter() + .flat_map(|row| row.error_metrics.keys().cloned()) + .collect::>(); + metrics + .into_iter() + .filter_map(|metric| { + let ErpParameterDecision::Empirical { + record_id, + params: selected, + .. + } = self.select_metric( + algorithm.clone(), + &metric, + f64::MAX, + params.clone(), + Some(params), + ) + else { + return None; + }; + if &selected != params { + return None; + } + let mut ids = if self.artifact.records.iter().any(|row| row.id == record_id) { + vec![record_id] + } else { + serde_json::from_str::>(&record_id).ok()? + }; + if ids.is_empty() { + return None; + } + // The logical proxy is per partition. For population-specific + // profiles use the largest matched partition, not a pooled sketch. + let mut resources = asap_aware_mapping::erp::ErpResourceProfile { + memory_bytes: 0.0, + update_cpu_seconds: 0.0, + merge_cpu_seconds: 0.0, + query_cpu_seconds: 0.0, + }; + for id in &ids { + let row = self.artifact.records.iter().find(|row| &row.id == id)?; + resources.memory_bytes = resources.memory_bytes.max(row.resources.memory_bytes); + resources.update_cpu_seconds = resources + .update_cpu_seconds + .max(row.resources.update_cpu_seconds); + resources.merge_cpu_seconds = resources + .merge_cpu_seconds + .max(row.resources.merge_cpu_seconds); + resources.query_cpu_seconds = resources + .query_cpu_seconds + .max(row.resources.query_cpu_seconds); + } + ids.sort(); + ids.dedup(); + Some((resources, ids)) + }) + .min_by(|a, b| { + a.0.memory_bytes + .total_cmp(&b.0.memory_bytes) + .then_with(|| a.1.cmp(&b.1)) + }) + } + + pub(crate) fn candidate_memory_bytes( + &self, + algorithm: &SketchAlgorithm, + params: &SketchParams, + ) -> Option<(f64, Vec)> { + self.candidate_resources(algorithm, params) + .map(|(resources, ids)| (resources.memory_bytes, ids)) + } + pub fn select( &self, algorithm: SketchAlgorithm, @@ -873,7 +1006,10 @@ fn u32_param(parameters: &Value, names: &[&str]) -> Option { .then_some(value as u32) } -fn parse_params(algorithm: &SketchAlgorithm, parameters: &Value) -> Option { +pub(crate) fn parse_params( + algorithm: &SketchAlgorithm, + parameters: &Value, +) -> Option { let width = || u32_param(parameters, &["width", "cols"]); let depth = || u32_param(parameters, &["depth", "rows"]); Some(match algorithm { @@ -1080,6 +1216,7 @@ mod tests { GroupingStrategy::PerSubpopulationInstance, ); let model = ErpAccuracyModel { + hll: None, policy: Some(&policy), max_error: 0.2, }; @@ -1093,6 +1230,7 @@ mod tests { )); policy.artifact.records[1].error_metrics.clear(); assert!(ErpAccuracyModel { + hll: None, policy: Some(&policy), max_error: 0.2 } @@ -1114,6 +1252,7 @@ mod tests { GroupingStrategy::PerSubpopulationInstance, ); let model = ErpAccuracyModel { + hll: None, policy: Some(&policy), max_error: 0.05, }; @@ -1160,6 +1299,7 @@ mod tests { GroupingStrategy::PerSubpopulationInstance, ); let model = ErpAccuracyModel { + hll: None, policy: Some(&policy), max_error: 0.2, }; @@ -1191,6 +1331,7 @@ mod tests { .error_metrics .remove("max_frequency_entropy_absolute_bits_error"); let model = ErpAccuracyModel { + hll: None, policy: Some(&policy), max_error: 0.2, }; @@ -1444,6 +1585,15 @@ mod tests { assert!(invalid.populations.is_empty()); assert!(invalid.invalid_reason.as_deref().unwrap().contains("stale")); assert!(policy.observed_shape.is_none()); + assert!(policy + .candidate_memory_bytes( + &SketchAlgorithm::Cms, + &SketchParams::Cms { + width: 512, + depth: 3 + } + ) + .is_none()); } /// Only the activated catalog may supply a candidate's data contract. @@ -1604,6 +1754,16 @@ mod tests { minimum_confidence: 0.9, minimum_confidence_margin: 0.05, }); + assert_eq!( + policy.candidate_memory_bytes( + &SketchAlgorithm::Cms, + &SketchParams::Cms { + width: 512, + depth: 3 + } + ), + Some((12_288.0, vec!["cms-512".into()])) + ); assert!(matches!( policy.select( SketchAlgorithm::Cms, @@ -1647,6 +1807,15 @@ mod tests { policy.select(SketchAlgorithm::Cms, 0.01, theory.clone()), ErpParameterDecision::TheoreticalFallback { params, .. } if params == theory )); + assert!(policy + .candidate_memory_bytes( + &SketchAlgorithm::Cms, + &SketchParams::Cms { + width: 512, + depth: 3 + } + ) + .is_none()); } #[test] diff --git a/control_plane/src/physical/post_asap/cost_model.rs b/control_plane/src/physical/post_asap/cost_model.rs index 9610ed0a1..8e0a4a3b5 100644 --- a/control_plane/src/physical/post_asap/cost_model.rs +++ b/control_plane/src/physical/post_asap/cost_model.rs @@ -167,6 +167,8 @@ pub struct ControlPlaneCostModel { offline_frequency_comparison: Option<(OfflineComparisonEvidence, OfflineComparisonRequest)>, exact_composition_costs: Vec, erp: Option, + erp_costs: Option, + hll_confidence: Option, } impl ControlPlaneCostModel { @@ -193,6 +195,8 @@ impl ControlPlaneCostModel { let dag = planner_types::post_asap::compile_executable_dag(root).ok()?; let mut value = 0.0; let mut states = 0; + let mut analytical = false; + let mut erp_record_ids = Vec::new(); for node in &dag.nodes { let family = match &node.payload { ExecutableOperatorPayload::SummaryAgg { family, .. } @@ -200,17 +204,38 @@ impl ControlPlaneCostModel { _ => continue, }; states += 1; - value += analytical_state_bytes(family)?; + let measurement = match family { + SummaryFamilyType::Sketch(kind, GroupingStrategy::PerSubpopulationInstance) => self + .erp_costs + .as_ref() + .and_then(|erp| erp.candidate_memory_bytes(kind.algorithm(), kind.params())), + _ => None, + }; + if let Some((bytes, ids)) = measurement { + value += bytes; + erp_record_ids.extend(ids); + } else { + value += analytical_state_bytes(family)?; + analytical = true; + } } if states == 0 || !value.is_finite() { return None; } + erp_record_ids.sort(); + erp_record_ids.dedup(); Some(CandidateCostEstimate { value, unit: "bytes_per_state_partition", model: "backend_state_footprint_v1", - source: "analytical", - erp_record_ids: vec![], + source: if erp_record_ids.is_empty() { + "analytical" + } else if analytical { + "mixed" + } else { + "erp" + }, + erp_record_ids, }) } @@ -224,6 +249,8 @@ impl ControlPlaneCostModel { offline_frequency_comparison: None, exact_composition_costs: Vec::new(), erp: None, + erp_costs: None, + hll_confidence: None, } } @@ -235,11 +262,27 @@ impl ControlPlaneCostModel { self } + pub fn with_hll_confidence( + mut self, + contract: super::super::erp::HllConfidenceContract, + ) -> Self { + self.hll_confidence = Some(contract); + self + } + pub fn with_erp(mut self, erp: ErpPlanningInput) -> Self { + self.erp_costs = Some(erp.clone()); self.erp = Some(erp); self } + /// Resource evidence can remain usable when ERP's error observations + /// cannot establish the query's requested confidence guarantee. + pub fn with_erp_costs(mut self, erp: ErpPlanningInput) -> Self { + self.erp_costs = Some(erp); + self + } + pub fn erp_parameter_decision( &self, algorithm: SketchAlgorithm, @@ -384,6 +427,44 @@ impl ControlPlaneCostModel { costs.into_iter().map(|(algorithm, _)| algorithm).collect() } + fn rank_with_erp_costs( + &self, + intent: &AggIntent, + defaults: Vec, + ) -> Vec { + let Some(erp) = &self.erp_costs else { + return defaults; + }; + let (eps, delta) = + asap_aware_mapping::replacement::accuracy_budget(&intent_accuracy(intent)); + let mut has_measurement = false; + let mut costs = defaults + .iter() + .map(|algorithm| { + let params = self.size_params(algorithm.clone(), intent, eps, delta); + let measured = erp.candidate_memory_bytes(algorithm, ¶ms); + has_measurement |= measured.is_some(); + let cost = measured.map(|(bytes, _)| bytes).or_else(|| { + analytical_state_bytes(&SummaryFamilyType::Sketch( + planner_types::post_asap::SketchKind::new(algorithm.clone(), params), + GroupingStrategy::PerSubpopulationInstance, + )) + }); + (algorithm.clone(), cost) + }) + .collect::>(); + if !has_measurement { + return defaults; + } + costs.sort_by(|(_, a), (_, b)| match (a, b) { + (Some(a), Some(b)) => a.total_cmp(b), + (Some(_), None) => std::cmp::Ordering::Less, + (None, Some(_)) => std::cmp::Ordering::Greater, + (None, None) => std::cmp::Ordering::Equal, + }); + costs.into_iter().map(|(algorithm, _)| algorithm).collect() + } + /// Keep concrete physical identities in Planner's complete-candidate estimate, /// including distinct pane sizes using the same abstract window framework. pub fn with_window_implementation_costs( @@ -693,7 +774,7 @@ impl CostModel for ControlPlaneCostModel { // puts it first (`summary_candidates`), nothing to reorder. _ => candidates.to_vec(), }; - self.rank_with_offline_evidence(intent, defaults) + self.rank_with_erp_costs(intent, self.rank_with_offline_evidence(intent, defaults)) } fn size_params( @@ -703,6 +784,19 @@ impl CostModel for ControlPlaneCostModel { eps: f64, delta: f64, ) -> SketchParams { + if kind == SketchAlgorithm::Hll && matches!(intent, AggIntent::Cardinality { .. }) { + if let Some(contract) = &self.hll_confidence { + if let Some((eps, delta)) = self.combined_eps_delta(&intent_accuracy(intent)) { + // If no supported precision meets the target, keep the tightest + // candidate for explain; its guarantee still fails selection. + let precision = contract + .confidence(eps) + .and_then(|model| model.precision(delta)) + .unwrap_or(18); + return SketchParams::Hll { precision }; + } + } + } let (max_error, theoretical) = match intent { AggIntent::TopK { k, .. } => { let (eps, delta) = self.topk_eps_delta(&intent_accuracy(intent)); @@ -998,7 +1092,7 @@ fn hll_precision_for_eps(eps: f64) -> u8 { } #[cfg(test)] -mod tests { +pub(crate) mod tests { use super::*; use planner_types::pre_asap::{default_cardinality, default_quantile}; @@ -1006,6 +1100,153 @@ mod tests { AccuracyTarget::Epsilon(e) } + /// Logical summary selection has an explicit Planner estimate before + /// deployment pricing, including when a family is forced. + #[test] + fn summary_candidates_have_explicit_logical_costs() { + use asap_aware_mapping::{ReplacementStrategy, SketchAlgorithmStrategy, TargetSubDAG}; + use std::rc::Rc; + + let accuracy = eps(0.1); + let root = Rc::new( + crate::query_parser::parse_query_expr_canonical( + "quantile_over_time(0.9,m[1m])", + accuracy.clone(), + ) + .unwrap(), + ); + let target = TargetSubDAG::new(&root); + let model = ControlPlaneCostModel::new(accuracy.clone()); + let forced = ForcedFamilyCostModel::new(accuracy, SketchAlgorithm::DDSketch); + let candidates = SketchAlgorithmStrategy::new(&model).replacements(&target); + assert!(!candidates.is_empty()); + for candidate in &candidates { + let estimate = model.candidate_cost_estimate(candidate).unwrap(); + assert!(estimate.value.is_finite() && estimate.value > 0.0); + assert_eq!(estimate.source, "analytical"); + assert_eq!(estimate.unit, "bytes_per_state_partition"); + assert_eq!( + model.candidate_cost(candidate, &target), + Some(Cost(estimate.value)) + ); + assert_eq!( + forced.candidate_cost(candidate, &target), + Some(Cost(estimate.value)) + ); + } + } + + /// Analytical estimates scale with configured state size; unsupported + /// families remain unavailable instead of receiving a made-up zero. + #[test] + fn analytical_footprint_scales_with_parameters() { + use planner_types::post_asap::SketchKind; + let kll = |k| { + SummaryFamilyType::Sketch( + SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k }), + GroupingStrategy::PerSubpopulationInstance, + ) + }; + assert_eq!(analytical_state_bytes(&kll(200)), Some(6400.0)); + assert_eq!(analytical_state_bytes(&kll(400)), Some(12800.0)); + let unsupported = SummaryFamilyType::Sketch( + SketchKind::new(SketchAlgorithm::Theta, SketchParams::Theta { k: 1024 }), + GroupingStrategy::PerSubpopulationInstance, + ); + assert_eq!(analytical_state_bytes(&unsupported), None); + } + + pub(crate) fn erp_cost_fixture() -> ErpPlanningInput { + serde_json::from_value(serde_json::json!({ + "artifact": {"schema_version": 1, "producer_version": "synthetic-test-fixture", + "records": [{"id": "kll-200", "sketch": "kll-percall", "implementation": "lib", + "parameters": {"k": 200}, "distribution": {"test": "population-a"}, "trials": 20, + "error_metrics": {"max_rank_err": 0.9}, + "resources": {"memory_bytes": 7777.0, "update_cpu_seconds": 1e-7, + "merge_cpu_seconds": 1e-5, "query_cpu_seconds": 1e-6}}]}, + "distribution": {"test": "population-a"}, "implementation": "lib", + "error_metric": "max_rank_err", "min_trials": 10, + "expected_updates": 1000.0, "expected_queries": 100.0, "expected_merges": 1.0, + "retention_seconds": 60.0, "cpu_weight": 1.0, "byte_second_weight": 1e-9, + "mode": "hybrid" + })).unwrap() + } + + /// ERP memory overrides analytical memory for the exact configuration, + /// even when it is larger. Mismatched profiles fall back to analytical. + #[test] + fn erp_cost_precedes_analytical_without_certifying_accuracy() { + use asap_aware_mapping::{ReplacementStrategy, SketchAlgorithmStrategy}; + use std::rc::Rc; + let policy = erp_cost_fixture(); + let root = Rc::new( + crate::query_parser::parse_query_expr_canonical( + "quantile_over_time(0.9,m[1m])", + eps(0.1), + ) + .unwrap(), + ); + let model = ControlPlaneCostModel::new(eps(0.1)); + let candidates = + SketchAlgorithmStrategy::new(&model).replacements(&TargetSubDAG::new(&root)); + let candidate = candidates + .iter() + .find(|candidate| { + model + .candidate_cost_estimate(candidate) + .is_some_and(|cost| cost.value == 6400.0) + }) + .expect("KLL candidate"); + let estimate = |policy: ErpPlanningInput| { + ControlPlaneCostModel::new(eps(0.1)) + .with_erp_costs(policy) + .candidate_cost_estimate(candidate) + .unwrap() + }; + let matched = estimate(policy.clone()); + assert_eq!(matched.value, 7777.0); + assert_eq!(matched.source, "erp"); + assert_eq!(matched.erp_record_ids, ["kll-200"]); + let mut mismatch = policy.clone(); + mismatch.distribution = serde_json::json!({"test": "population-b"}); + assert_eq!(estimate(mismatch).source, "analytical"); + let mut mismatch = policy.clone(); + mismatch.implementation = Some("different-runtime".into()); + assert_eq!(estimate(mismatch).source, "analytical"); + let mut mismatch = policy.clone(); + mismatch.artifact.records[0].parameters = serde_json::json!({"k": 400}); + assert_eq!(estimate(mismatch).source, "analytical"); + let mut mismatch = policy.clone(); + mismatch.min_trials = 21; + assert_eq!(estimate(mismatch).source, "analytical"); + let mut mismatch = policy.clone(); + mismatch.artifact.records[0].resources.memory_bytes = f64::NAN; + assert_eq!(estimate(mismatch).source, "analytical"); + let model = ControlPlaneCostModel::new(eps(0.1)).with_erp_costs(policy); + assert_eq!( + model.rank_candidates( + &default_quantile(0.9), + &[SketchAlgorithm::DDSketch, SketchAlgorithm::Kll] + )[0], + SketchAlgorithm::Kll + ); + let (_, trace) = crate::planner_selection::select_workload_with_accuracy_model_and_trace( + vec![(0, root)], + eps(0.1), + &model, + &asap_aware_mapping::NoAccuracyEvidence, + &asap_aware_mapping::DefaultAccuracyModel, + ) + .unwrap(); + assert!(trace["groups"] + .as_array() + .unwrap() + .iter() + .flat_map(|group| group["candidates"].as_array().unwrap()) + .any(|candidate| candidate["cost_estimate"]["source"] == "erp" + && candidate["cost_estimate"]["unit"] == "bytes_per_state_partition")); + } + #[test] fn quantile_always_prefers_ddsketch_over_kll() { let model = ControlPlaneCostModel::new(AccuracyTarget::Epsilon(0.1)); diff --git a/control_plane/src/physical/workload_cost.rs b/control_plane/src/physical/workload_cost.rs index 5a0cef586..eea4d3532 100644 --- a/control_plane/src/physical/workload_cost.rs +++ b/control_plane/src/physical/workload_cost.rs @@ -4,7 +4,9 @@ //! Planner owns semantic legality. A manifest describes the exact physical //! demand to price; provider quotes and candidate evaluations are separate. +mod automatic; mod materialization_candidates; +pub use automatic::{AutomaticCostReport, ComponentResources}; mod status; pub use status::{CandidateEvaluationStatus, CandidateSearchScope}; @@ -91,6 +93,8 @@ pub struct CandidatePlanEvaluation { pub status: CandidateEvaluationStatus, pub plan_id: Option, pub total_cost: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub automatic_cost: Option, pub unavailable_reason: Option, } @@ -529,6 +533,7 @@ fn candidate_description(candidate: &PhysicalCompilationRequest) -> CandidatePla status: CandidateEvaluationStatus::CompilationFailed, plan_id: None, total_cost: None, + automatic_cost: None, unavailable_reason: None, } } @@ -546,6 +551,19 @@ fn compile_candidate_for_pricing( Box, > { let mut description = candidate_description(&candidate); + if let Some(policy) = &candidate.erp { + if candidate.queries.iter().any(|query| { + super::compiler::requires_exact_erp_fallback( + &query.selected_plan_root, + &query.accuracy_target, + policy, + ) + }) { + description.unavailable_reason = + Some("candidate does not satisfy ERP accuracy admission".into()); + return Err(Box::new(description)); + } + } let queries = candidate.queries.clone(); let compiled = super::realization::RealizationProvider::compile( &super::realization::ExistingRealizations, @@ -648,7 +666,7 @@ pub fn select_lowest_cost_candidate( select_candidates( candidates, env, - evidence, + Some(evidence), super::compiler::QueryFrontend::PromQl, ) } @@ -661,18 +679,20 @@ pub fn select_lowest_cost_metricsql_candidate( select_candidates( candidates, env, - evidence, + Some(evidence), super::compiler::QueryFrontend::MetricsQl, ) } -fn select_candidates( +pub fn select_candidates( candidates: Vec, env: PhysicalDeploymentContext, - evidence: &WorkloadCostEvidence, + evidence: Option<&WorkloadCostEvidence>, frontend: super::compiler::QueryFrontend, ) -> Result { - evidence.validate(&env)?; + if let Some(evidence) = evidence { + evidence.validate(&env)?; + } if candidates.is_empty() || candidates.len() > 4096 { return Err(invalid( "candidate inventory must contain 1..=4096 candidates", @@ -697,9 +717,23 @@ fn select_candidates( }); let planner_selection_trace = candidates[0].planner_selection_trace.clone(); let mut comparison_workload = None; + let mut comparison_inputs = None; let mut candidate_evaluations = Vec::new(); let mut priced_candidates = Vec::new(); for candidate in candidates { + if evidence.is_none() { + let inputs = json!({"data":candidate.data_workload,"erp":candidate.erp}); + if comparison_inputs + .as_ref() + .is_some_and(|previous| previous != &inputs) + { + return Err(invalid( + "candidates describe different data/ERP cost inputs", + )); + } + comparison_inputs = Some(inputs); + } + let cost_request = candidate.clone(); let (plan, manifest, mut description) = match compile_candidate_for_pricing(candidate, env.clone(), frontend) { Ok(bound) => bound, @@ -716,11 +750,32 @@ fn select_candidates( return Err(invalid("candidates describe different workloads/horizons")); } comparison_workload = Some(scope); - match super::realization::RealizationProvider::price( - &super::realization::ExistingRealizations, - evidence, - &manifest, - ) { + let priced = if let Some(evidence) = evidence { + super::realization::RealizationProvider::price( + &super::realization::ExistingRealizations, + evidence, + &manifest, + ) + } else { + automatic::estimate(&cost_request, &env, &plan, &manifest) + .map(|report| { + let costs = report + .components + .iter() + .map(|(id, r)| (id.clone(), r.weighted_cost())) + .collect::>(); + let total = Cost(costs.values().sum()); + description.automatic_cost = Some(report); + (total, costs) + }) + .map_err(|error| { + ( + CandidateEvaluationStatus::EvidenceMissing, + error.to_string(), + ) + }) + }; + match priced { Ok((cost, components)) => { description.status = CandidateEvaluationStatus::Unselected; description.total_cost = Some(cost.0); @@ -765,8 +820,12 @@ fn select_candidates( plan.cost_comparison = Some(CandidatePlanSelectionReport { planner_selection_trace, materialization_search_coverage, - data_snapshot_id: evidence.data_snapshot_id.clone(), - model_version: evidence.model_version.clone(), + data_snapshot_id: evidence + .map(|e| e.data_snapshot_id.clone()) + .unwrap_or_else(|| "request-scoped-data-workload".into()), + model_version: evidence + .map(|e| e.model_version.clone()) + .unwrap_or_else(|| automatic::MODEL_VERSION.into()), selected_plan_id: plan.envelope.plan_id, selected_manifest, component_costs, @@ -1014,7 +1073,7 @@ mod tests { .contains("external execution is unavailable"), "{error}" ); - let selected = with_unit_quotes(input).compile_promql().unwrap(); + let selected = input.compile_promql().unwrap(); assert!(selected.query_plan.entries.values().all(|entry| entry .nodes .values() @@ -1072,7 +1131,7 @@ mod tests { let queries = input.query_workload.repeating_queries.as_mut().unwrap(); queries.truncate(1); queries[0].query = planner_types::workload::Query("count by(job)(m)".into()); - let plan = with_unit_quotes(input).compile_promql().unwrap(); + let plan = input.compile_promql().unwrap(); assert!(plan .query_plan .entries @@ -1089,6 +1148,153 @@ mod tests { )))); } + /// Deployment computes and compares complete costs without external quotes. + #[test] + fn deployment_automatically_prices_workload() { + let mut input = fixture(); + input.workload_cost_evidence = None; + let plan = input.compile_promql().unwrap(); + let report = plan.cost_comparison.unwrap(); + let selected = report + .candidate_evaluations + .iter() + .find(|c| c.status == CandidateEvaluationStatus::Selected) + .unwrap(); + assert_eq!( + selected.total_cost.unwrap(), + report.component_costs.values().sum::() + ); + assert!(report + .candidate_evaluations + .iter() + .filter_map(|c| c.total_cost) + .all(|cost| cost >= selected.total_cost.unwrap())); + assert!(!report.component_costs.is_empty()); + assert_eq!( + report.component_costs.len(), + report.selected_manifest.components.len() + ); + assert!( + report + .candidate_evaluations + .iter() + .filter(|c| c.total_cost.is_some()) + .count() + > 1 + ); + } + + fn automatic_report( + request: &PhysicalCompilationRequest, + env: &PhysicalDeploymentContext, + plan: &CompiledPhysicalPlan, + ) -> AutomaticCostReport { + automatic::estimate( + request, + env, + plan, + &manifest(plan, &request.queries).unwrap(), + ) + .unwrap() + } + + /// Demand changes recurring work, while shared maintenance stays once per location. + #[test] + fn automatic_cost_scales_demand_data_and_shared_consumers() { + let (mut request, env) = fixture().into_physical_compilation_request().unwrap(); + let plan = DeploymentPlanCompiler + .compile_promql(request.clone(), env.clone()) + .unwrap(); + let before = automatic_report(&request, &env, &plan); + request.queries[0] + .summary_lifecycle_inputs + .evaluation_interval_ms /= 2; + let frequent = automatic_report(&request, &env, &plan); + for (id, cost) in &before.components { + let factor = if id.starts_with("query:") || id.starts_with("result:") { + 2.0 + } else { + 1.0 + }; + assert!( + (frequent.components[id].weighted_cost() - factor * cost.weighted_cost()).abs() + < 1e-10, + "{id}" + ); + } + request + .data_workload + .as_mut() + .unwrap() + .ingestion_rate + .value + .as_mut() + .unwrap() + .0 *= 2.0; + let larger = automatic_report(&request, &env, &plan); + assert!( + larger + .components + .iter() + .filter(|(id, _)| id.ends_with(":update")) + .map(|(_, c)| c.cpu_seconds) + .sum::() + > frequent + .components + .iter() + .filter(|(id, _)| id.ends_with(":update")) + .map(|(_, c)| c.cpu_seconds) + .sum::() + ); + let mut shared = plan.clone(); + let mut query = request.queries[0].clone(); + query.query_id = "second-consumer".into(); + let mut entry = shared.query_plan.entries.values().next().unwrap().clone(); + entry.query_id = query.query_id.clone(); + shared + .query_plan + .entries + .insert(query.query_id.clone(), entry); + request.queries.push(query); + let twice = automatic_report(&request, &env, &shared); + for (id, cost) in &larger.components { + assert_eq!(&twice.components[id], cost); + } + assert_eq!( + twice + .components + .keys() + .filter(|id| id.starts_with("source:") || id.starts_with("state:")) + .count(), + larger + .components + .keys() + .filter(|id| id.starts_with("source:") || id.starts_with("state:")) + .count() + ); + } + + /// Expired facts and arithmetic overflow cannot silently become a zero quote. + #[test] + fn automatic_cost_rejects_stale_and_overflowing_data() { + let (mut request, env) = fixture().into_physical_compilation_request().unwrap(); + let plan = DeploymentPlanCompiler + .compile_promql(request.clone(), env.clone()) + .unwrap(); + let manifest = manifest(&plan, &request.queries).unwrap(); + let rate = &mut request.data_workload.as_mut().unwrap().ingestion_rate; + rate.observed_at_ms = Some(1); + rate.valid_for_ms = Some(1); + assert!(automatic::estimate(&request, &env, &plan, &manifest) + .unwrap_err() + .to_string() + .contains("fresh ingestion rate")); + let rate = &mut request.data_workload.as_mut().unwrap().ingestion_rate; + rate.valid_for_ms = None; + rate.value.as_mut().unwrap().0 = f64::MAX; + assert!(automatic::estimate(&request, &env, &plan, &manifest).is_err()); + } + fn fixture() -> BackendLocalPlanningInput { let mut snapshot: BackendLocalPlanningInput = serde_json::from_str(include_str!( "../../../docs/examples/asapquery-planning-snapshot.json" @@ -1395,45 +1601,6 @@ mod tests { (candidates, env, evidence) } - /// Quote every bindable candidate at unit cost, leaving admission to decide. - fn with_unit_quotes(mut input: BackendLocalPlanningInput) -> BackendLocalPlanningInput { - let (request, env) = input.clone().into_physical_compilation_request().unwrap(); - let quotes = enumerate_exact_and_materialized_candidates(request) - .unwrap() - .into_iter() - .filter_map(|candidate| { - let plan = DeploymentPlanCompiler - .compile_promql(candidate.clone(), env.clone()) - .ok()?; - let manifest = manifest(&plan, &candidate.queries).unwrap(); - let unit_costs = manifest - .components - .keys() - .map(|id| (id.clone(), 1.0)) - .collect(); - Some(WorkloadQuote { - manifest, - executable: true, - unit_costs, - }) - }) - .collect(); - input.workload_cost_evidence = Some(WorkloadCostEvidence { - backend_revision: crate::physical::compiler::BACKEND_REVISION.into(), - planner_revision: crate::physical::compiler::PLANNER_REVISION.into(), - data_snapshot_id: input - .physical_inputs - .data_snapshot_id - .clone() - .unwrap_or_else(|| "fixture-data-v1".into()), - model_version: "test-only-unit-costs".into(), - observed_at_unix_ms: env.observed_at_unix_ms, - valid_for_ms: env.max_evidence_age_ms, - quotes, - }); - input - } - // Retained local input is priced once per metric, separate from the native service. #[test] fn counter_materialization_manifest_prices_owned_state_and_distinct_native_candidate() { @@ -1640,9 +1807,14 @@ mod tests { } #[test] - fn snapshot_requires_quotes_and_roundtrips_selection() { + fn snapshot_supports_automatic_cost_and_roundtrips_quote_override() { let mut snapshot = fixture(); - assert!(snapshot.clone().compile_promql().is_err()); + assert!(snapshot + .clone() + .compile_promql() + .unwrap() + .cost_comparison + .is_some()); let (_, _, evidence) = quoted(); snapshot.workload_cost_evidence = Some(evidence); let snapshot: BackendLocalPlanningInput = diff --git a/control_plane/src/physical/workload_cost/automatic.rs b/control_plane/src/physical/workload_cost/automatic.rs new file mode 100644 index 000000000..6f0be5883 --- /dev/null +++ b/control_plane/src/physical/workload_cost/automatic.rs @@ -0,0 +1,810 @@ +//! Versioned reference resource model. These coefficients are analytical +//! assumptions, not measurements or currency prices. ERP replaces only the +//! resource dimensions it measures; plan demand always comes from the backend. +use super::*; +use crate::query_plan::{residual::ResidualQueryOperator as Op, QueryPlanNode as Node}; +use asap_aware_mapping::erp::ErpResourceProfile; +use asap_types::{ + AggregationType as A, PrecomputeMaterialization, WindowMaterializationLayout as Layout, +}; +use planner_types::post_asap::SketchAlgorithm; + +pub const MODEL_VERSION: &str = "backend-workload-resources-v2"; +const CPU_PER_ITEM: f64 = 1e-7; +const CPU_PER_BYTE: f64 = 1e-9; +const SAMPLE_BYTES: f64 = 24.0; +const SERIES_BYTES: f64 = 256.0; +const MEMORY_WEIGHT: f64 = 1e-9; +const NETWORK_WEIGHT: f64 = 1e-8; + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct ComponentResources { + pub cpu_seconds: f64, + pub memory_byte_seconds: f64, + pub network_bytes: f64, + pub source: String, + pub erp_record_ids: Vec, + pub calculation: Value, +} +impl ComponentResources { + pub(super) fn weighted_cost(&self) -> f64 { + self.cpu_seconds + + self.memory_byte_seconds * MEMORY_WEIGHT + + self.network_bytes * NETWORK_WEIGHT + } +} +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct AutomaticCostReport { + pub model_version: String, + pub weights: Value, + pub inputs: Value, + pub assumptions: Vec, + pub components: BTreeMap, +} +#[derive(Clone)] +struct State { + profile: ErpResourceProfile, + ids: Vec, + partitions: f64, + retained: f64, + allocations_per_second: f64, + overlap: f64, + rollup_merges: f64, +} +fn resources( + cpu: f64, + memory: f64, + network: f64, + ids: Vec, + calculation: Value, +) -> ComponentResources { + ComponentResources { + cpu_seconds: cpu, + memory_byte_seconds: memory, + network_bytes: network, + source: if ids.is_empty() { + "analytical" + } else { + "erp+analytical" + } + .into(), + erp_record_ids: ids, + calculation, + } +} +// Price the installed program, including transient heap state. Group count is +// conservatively bounded by input cardinality until scoped group statistics exist. +fn snapshot_work( + entry: &asap_types::query_plan::QueryPlanEntry, + rows: f64, +) -> Result<(f64, f64), CompileError> { + let program = entry + .physical_dag + .as_ref() + .ok_or_else(|| invalid("missing snapshot physical program"))?; + physical_work(program, rows) +} +fn physical_work(program: &Value, rows: f64) -> Result<(f64, f64), CompileError> { + use planner_types::post_asap::{SketchParams, SummaryFamilyType}; + let nodes = program["nodes"] + .as_object() + .ok_or_else(|| invalid("invalid snapshot program"))?; + let mut items = rows; // Bind complete input rows. + let mut bytes = rows * SERIES_BYTES * 3.0; + for node in nodes.values().filter_map(|node| node.get("Operator")) { + let kind = node["operator"]["kind"] + .as_object() + .ok_or_else(|| invalid("invalid native operator"))?; + let Some((name, parameters)) = kind.iter().next() else { + return Err(invalid("missing native operator kind")); + }; + match name.as_str() { + "Sort" => items += rows * rows.max(2.0).log2(), + "Limit" | "Project" | "CompiledProject" => items += rows, + "KeyedSummaryBuild" => { + let family: SummaryFamilyType = + serde_json::from_value(parameters["family"].clone()) + .map_err(|error| invalid(error.to_string()))?; + let SummaryFamilyType::Sketch(kind, _) = family else { + return Err(invalid("no snapshot cost model for non-sketch keyed build")); + }; + let (width, depth, capacity) = match kind.params() { + SketchParams::CmsWithHeap { + width, + depth, + heap_size, + } + | SketchParams::CountSketchWithHeap { + width, + depth, + heap_size, + } => (*width as f64, *depth as f64, *heap_size as f64), + _ => return Err(invalid("no snapshot heap cost model for this algorithm")), + }; + let groups = if parameters["groups"] + .as_array() + .is_some_and(|groups| groups.is_empty()) + { + rows.min(1.0) + } else { + rows + }; + items += rows * (depth + capacity.max(2.0).log2()); + bytes += groups * (width * depth * 8.0 + capacity * (SERIES_BYTES + 24.0)); + } + "SummaryBuild" => { + let family: SummaryFamilyType = + serde_json::from_value(parameters["family"].clone()) + .map_err(|e| invalid(e.to_string()))?; + if !matches!( + family, + SummaryFamilyType::ExactAggregate(planner_types::post_asap::ExactKind::Sum, _) + ) { + return Err(invalid("no native aggregate cost for this family")); + } + items += rows; + bytes += rows * 64.0; + } + "Readout" => items += rows, + "KeyedReadout" => items += rows * rows.max(2.0).log2(), + _ => { + return Err(invalid(format!( + "no automatic snapshot cost model for {name}" + ))) + } + } + } + Ok((items * CPU_PER_ITEM, bytes)) +} + +fn state( + m: &PrecomputeMaterialization, + request: &PhysicalCompilationRequest, + cardinality: u64, +) -> Result { + let bytes = + super::super::compiler::estimated_state_bytes(&m.aggregation_type, &m.parameters) as f64; + let mut profile = ErpResourceProfile { + memory_bytes: bytes, + update_cpu_seconds: CPU_PER_ITEM * (bytes / 16.0).max(2.0).log2(), + merge_cpu_seconds: CPU_PER_BYTE * bytes, + query_cpu_seconds: CPU_PER_BYTE * bytes + CPU_PER_ITEM, + }; + let mut ids = Vec::new(); + let algorithm = match m.aggregation_type { + A::DatasketchesKLL => Some(SketchAlgorithm::Kll), + A::HLL => Some(SketchAlgorithm::Hll), + A::UnivMon => Some(SketchAlgorithm::UnivMon), + _ => None, + }; + if let (Some(policy), Some(algorithm)) = (&request.erp, algorithm) { + let mut policy = policy.clone(); + // Use the same deployed implementation and population checks as the + // logical cost path. Unknown implementations fall back to this model. + policy.artifact.records.retain(|row| { + (row.sketch == "kll-percall" && row.implementation == "lib") + || (row.sketch == "hll" && row.implementation == "asap-sketchlib-hll-regular-v1") + || (row.sketch == "univmon" + && row.implementation == "asap-sketchlib-univmon-standard-v1") + }); + let population_matches = policy.observed_populations.is_none() + || (!request.canonical_roots.is_empty() + && request.canonical_roots.iter().all(|root| { + super::super::compiler::observed_population_matches_root(&policy, root) + })); + if population_matches { + if let Some(params) = super::super::erp::parse_params(&algorithm, &json!(m.parameters)) + { + if let Some((measured, records)) = policy.candidate_resources(&algorithm, ¶ms) { + profile = measured; + ids = records; + } + } + } + } + let (period, overlap, rollup_merges) = match &m.window_layout { + Layout::Pane { pane_secs } => (*pane_secs as f64, 1.0, 0.0), + Layout::FullWindow => ( + m.slide_interval as f64, + (m.window_size as f64 / m.slide_interval as f64).ceil(), + 0.0, + ), + Layout::HierarchicalRollup { + base_pane_secs, + levels_secs, + } => { + // Every closed child state is folded into its parent once. + let periods = std::iter::once(base_pane_secs) + .chain(levels_secs.iter().take(levels_secs.len().saturating_sub(1))); + ( + *base_pane_secs as f64, + 1.0, + periods.map(|p| 1.0 / *p as f64).sum(), + ) + } + }; + let retained = m + .num_aggregates_to_retain + .ok_or_else(|| invalid("state retention count is unknown"))? as f64; + if period <= 0.0 || retained <= 0.0 { + return Err(invalid("invalid state layout")); + } + Ok(State { + profile, + ids, + partitions: super::super::compiler::retained_partition_count(m, Some(cardinality)) as f64, + retained, + allocations_per_second: 1.0 / period + + match &m.window_layout { + Layout::HierarchicalRollup { levels_secs, .. } => { + levels_secs.iter().map(|p| 1.0 / *p as f64).sum::() + } + _ => 0.0, + }, + overlap, + rollup_merges, + }) +} + +pub(super) fn estimate( + request: &PhysicalCompilationRequest, + env: &PhysicalDeploymentContext, + plan: &CompiledPhysicalPlan, + manifest: &WorkloadCostManifest, +) -> Result { + let data = request + .data_workload + .as_ref() + .ok_or_else(|| invalid("automatic costing needs data_workload"))?; + let now = env.observed_at_unix_ms; + let horizon = manifest.horizon_seconds; + let rate = data + .ingestion_rate + .value_at(now) + .map(|r| r.0) + .ok_or_else(|| invalid("automatic costing needs a fresh ingestion rate"))?; + let cadence = data + .data_ingestion_interval + .value_at(now) + .map(|d| d.0 as f64 / 1000.0) + .ok_or_else(|| invalid("automatic costing needs fresh source cadence"))?; + if !rate.is_finite() || rate < 0.0 || !cadence.is_finite() || cadence <= 0.0 { + return Err(invalid("invalid source rate/cadence")); + } + let mut assumptions = vec![ + "Reference analytical coefficients are estimates, not calibrated benchmarks or currency prices.".into(), + "Absent per-source/selectivity facts, each source and collector uses the whole workload rate and cardinality (conservative replication).".into(), + "Grouped states use input cardinality as an upper bound on group count; filtered inputs receive no selectivity discount.".into(), + "Transport uses full state size even for deltas; memory footprint approximates encoded size; result rows include a 256-byte label allowance.".into(), + ]; + let cardinality = match data.input_cardinality.value_at(now) { + Some(value) => *value, + None => { + if rate == 0.0 { + return Err(invalid( + "zero arrival rate cannot establish active cardinality", + )); + } + assumptions.push("Unknown cardinality estimated as ceil(ingestion_rate * source_cadence), assuming one sample per active series per scrape.".into()); + let n = (rate * cadence).ceil(); + if !n.is_finite() || n >= u64::MAX as f64 { + return Err(invalid("cardinality estimate overflow")); + } + n as u64 + } + }; + if cardinality == 0 && rate > 0.0 { + return Err(invalid("positive ingestion rate with zero cardinality")); + } + let cardinality = cardinality as f64; + let mut states = BTreeMap::new(); + for schema in &plan.precompute_plan.schemas { + let m = plan + .precompute_plan + .materializations + .iter() + .find(|m| m.policy_fingerprint() == schema.materialization.fingerprint()) + .ok_or_else(|| invalid("missing materialization"))?; + states.insert( + schema.materialization.fingerprint(), + state(m, request, cardinality as u64)?, + ); + } + let mut components = BTreeMap::new(); + let mut insert = |id: String, value: ComponentResources| -> Result<(), CompileError> { + if !manifest.components.contains_key(&id) { + return Err(invalid(format!("unmanifested cost {id}"))); + } + if [ + value.cpu_seconds, + value.memory_byte_seconds, + value.network_bytes, + value.weighted_cost(), + ] + .iter() + .any(|v| !v.is_finite() || *v < 0.0) + { + return Err(invalid(format!("non-finite/negative resource cost {id}"))); + } + if components.insert(id.clone(), value).is_some() { + return Err(invalid(format!("duplicate cost {id}"))); + } + Ok(()) + }; + let updates = rate * horizon; + let max_lookback = plan + .query_plan + .entries + .values() + .map(|e| e.instant.lookback_ms as f64 / 1000.0) + .fold(cadence, f64::max); + let needs_volume = matches!(data.arrival, planner_types::workload::DataArrival::AtRest) + || plan + .query_plan + .entries + .values() + .any(|e| e.instant.full_history); + let source_samples = + if needs_volume { + *data.ingestion_volume.value_at(now).ok_or_else(|| { + invalid("at-rest/full-history costing needs fresh ingestion_volume") + })? as f64 + + updates + } else { + rate * max_lookback + }; + for (id, demand) in &manifest.components { + if id.starts_with("source:") { + let location = demand.implementation["location"].as_str().unwrap_or(""); + let remote_backend = location == "backend" && !plan.collector_plans.is_empty(); + let items = if remote_backend { 0.0 } else { updates }; + let raw_retained = if location == "exact_backend" { + source_samples * SAMPLE_BYTES + cardinality * SERIES_BYTES + } else { + 0.0 + }; + insert( + id.clone(), + resources( + items * CPU_PER_ITEM, + raw_retained * horizon, + items * SAMPLE_BYTES, + vec![], + json!({"input_samples": items, "retained_raw_bytes": raw_retained, "cpu_seconds_per_sample": CPU_PER_ITEM, + "backend_summary_decode_charged_in_transport": remote_backend}), + ), + )?; + } else if id.starts_with("maintenance:") { + let binding = &demand.implementation; + let interval = binding["interval_ms"] + .as_u64() + .filter(|value| *value > 0) + .ok_or_else(|| invalid("invalid native maintenance interval"))? + as f64 + / 1000.0; + let window = binding["window_ms"] + .as_u64() + .ok_or_else(|| invalid("missing native window"))? as f64 + / 1000.0; + let (cpu, workspace) = physical_work(&binding["physical_program"], cardinality)?; + let runs = (horizon / interval).ceil(); + let counter_scan = cardinality * (window / cadence).ceil() * CPU_PER_ITEM; + insert( + id.clone(), + resources( + runs * (cpu + counter_scan), + workspace * horizon, + 0.0, + vec![], + json!({"runs":runs,"input_series":cardinality,"window_seconds":window,"workspace_bytes_bound":workspace,"physical_program":binding["physical_program"]}), + ), + )?; + } else if id.starts_with("state:") { + let binding = &demand.implementation["binding"]; + let materialization: asap_types::PolicyFingerprint = + serde_json::from_value(binding["schema"]["materialization"].clone()) + .map_err(|e| invalid(e.to_string()))?; + let s = states + .get(&materialization) + .ok_or_else(|| invalid("unknown state"))?; + let operation = demand.implementation["operation"].as_str().unwrap_or(""); + let builds = s.partitions * (s.retained + (horizon * s.allocations_per_second).ceil()); + let remote_backend = + binding["location"] == "backend" && !plan.collector_plans.is_empty(); + let update_count = if remote_backend { + 0.0 + } else { + updates * s.overlap + }; + let merges = s.partitions * horizon * s.rollup_merges; + let (cpu, memory) = match operation { + "build" => (builds * s.profile.memory_bytes * CPU_PER_BYTE, 0.0), + "update" => ( + update_count * s.profile.update_cpu_seconds + + merges * s.profile.merge_cpu_seconds, + 0.0, + ), + "residency" => ( + 0.0, + s.partitions * s.retained * s.profile.memory_bytes * horizon, + ), + "retire" => (builds * CPU_PER_ITEM, 0.0), + _ => return Err(invalid("unsupported state phase")), + }; + insert( + id.clone(), + resources( + cpu, + memory, + 0.0, + s.ids.clone(), + json!({"operation":operation, + "partitions":s.partitions, "retained_states_per_partition":s.retained, "builds_and_retirements":builds, + "updates":update_count,"rollup_merges":merges,"unit_resources":s.profile, + "allocation_cpu_seconds_per_byte":CPU_PER_BYTE,"retirement_cpu_seconds_per_state":CPU_PER_ITEM}), + ), + )?; + } else if id.starts_with("current-series:") { + let phase = demand.implementation["phase"].as_str().unwrap_or(""); + let bytes = cardinality * SERIES_BYTES; + let retention_ms = demand.implementation["population"]["history_retention_ms"] + .as_u64() + .unwrap_or(0); + let versions = if retention_ms == 0 { + 0.0 + } else { + (retention_ms as f64 / (cadence * 1000.)).ceil() + 1.0 + }; + let retained_bytes = bytes + versions * (bytes + 1024.0); + let snapshots = if versions == 0.0 { + 0.0 + } else { + horizon / cadence + }; + let (cpu, memory) = match phase { + "build" => (bytes * CPU_PER_BYTE, 0.0), + "update" => ( + updates * CPU_PER_ITEM * cardinality.max(2.0).log2() + + snapshots * bytes * CPU_PER_BYTE, + 0.0, + ), + "residency" => (0.0, retained_bytes * horizon), + "retire" => ((cardinality + snapshots) * CPU_PER_ITEM, 0.0), + _ => return Err(invalid("unsupported current-series phase")), + }; + insert( + id.clone(), + resources( + cpu, + memory, + 0.0, + vec![], + json!({"phase":phase,"series":cardinality,"updates":updates,"bytes_per_series":SERIES_BYTES,"historical_versions":versions,"retained_bytes":retained_bytes,"snapshots":snapshots}), + ), + )?; + } + } + for rule in &plan.transmission_plan.rules { + if rule.emit_every_ms == 0 { + return Err(invalid("unknown transmission cadence")); + } + let s = states + .get(&rule.materialization.fingerprint()) + .ok_or_else(|| invalid("unknown transport state"))?; + let checkpoints = rule + .full_checkpoint_every_ms + .map(|ms| { + if ms == 0 { + f64::INFINITY + } else { + horizon * 1000.0 / ms as f64 + } + }) + .unwrap_or(0.0); + let frames = (horizon * 1000.0 / rule.emit_every_ms as f64 + checkpoints).ceil() + * s.partitions + * s.retained; + let bytes = frames * s.profile.memory_bytes; + insert( + format!( + "transport:{}:{}", + rule.producer_id, + rule.materialization.fingerprint() + ), + resources( + bytes * CPU_PER_BYTE * 2.0 + frames * s.profile.merge_cpu_seconds, + 0.0, + bytes, + s.ids.clone(), + json!({"frames":frames,"bytes_per_frame":s.profile.memory_bytes,"merge_cpu_seconds_per_frame":s.profile.merge_cpu_seconds, + "encode_decode_cpu_seconds_per_byte":2.0*CPU_PER_BYTE}), + ), + )?; + } + for entry in plan.query_plan.entries.values() { + let mut rows = BTreeMap::new(); + let mut profiles: BTreeMap<_, Vec<&State>> = BTreeMap::new(); + for node_id in entry + .topological_order() + .map_err(|e| invalid(e.to_string()))? + { + let node = &entry.nodes[&node_id]; + let id = format!("query:{}:{}", entry.query_id, node_id.0); + let evaluations = manifest.components[&id].occurrences_per_horizon; + let input_rows: f64 = node.inputs().iter().map(|id| rows[id]).sum(); + let mut output_rows = input_rows.max(1.0); + let mut used: Vec<&State> = node + .inputs() + .iter() + .flat_map(|id| profiles.get(id).into_iter().flatten().copied()) + .collect(); + let mut detail = json!({"input_rows":input_rows}); + let mut network = 0.0; + let cpu = match node { + Node::PhysicalFragment { dag, .. } => { + let program: Value = + serde_json::from_slice(dag).map_err(|error| invalid(error.to_string()))?; + let (cpu, workspace) = physical_work(&program, input_rows)?; + detail = json!({"input_rows":input_rows,"physical_program":program,"workspace_bytes_bound":workspace}); + cpu + } + Node::Physical { .. } => { + let (cpu, workspace) = snapshot_work(entry, input_rows)?; + detail = json!({"input_rows":input_rows,"physical_program":entry.physical_dag,"workspace_bytes_bound":workspace}); + cpu + } + Node::ReadMaterialization { binding } => { + let s = states + .get(&binding.materialization.fingerprint()) + .ok_or_else(|| invalid("unknown read state"))?; + let panes = if binding.full_window_slide_ms.is_some() { + 1.0 + } else { + (binding + .readout_lookback_ms + .unwrap_or(entry.instant.lookback_ms) as f64 + / binding.window_ms as f64) + .ceil() + .max(1.0) + }; + output_rows = s.partitions; + used = vec![s]; + detail = json!({"partitions":s.partitions,"panes_per_read":panes,"unit_resources":s.profile}); + s.partitions + * (panes * s.profile.memory_bytes * CPU_PER_BYTE + + (panes - 1.0) * s.profile.merge_cpu_seconds) + } + Node::SummaryEstimate { .. } => used + .iter() + .map(|s| s.partitions * s.profile.query_cpu_seconds) + .sum(), + Node::SummaryMerge { .. } => used + .iter() + .map(|s| s.partitions * s.profile.merge_cpu_seconds) + .sum(), + Node::Scalar { .. } => { + output_rows = 1.0; + CPU_PER_ITEM + } + Node::ExactFallback { .. } + | Node::Logical { + operator: Op::ExactSubquery { .. } | Op::CandidateExactSubquery { .. }, + .. + } => { + let query = match node { + Node::Logical { + operator: + Op::ExactSubquery { query } | Op::CandidateExactSubquery { query, .. }, + .. + } => query.as_str(), + _ => entry.canonical_query.as_str(), + }; + let parsed = crate::query_parser::parse_query_expr_canonical( + query, + crate::types::AccuracyTarget::Exact, + ) + .map_err(|e| invalid(e.to_string()))?; + let sources = exact_source_metrics(&parsed)?.len() as f64; + let samples = if needs_volume { + *data.ingestion_volume.value_at(now).ok_or_else(|| { + invalid("full-history exact cost needs fresh ingestion_volume") + })? as f64 + } else { + (rate * (entry.instant.lookback_ms as f64 / 1000.0).max(cadence)) + .max(cardinality) + } * sources; + let operators = tree_size( + &serde_json::to_value(parsed).map_err(|e| invalid(e.to_string()))?, + ) as f64; + output_rows = cardinality.max(1.0); + network = output_rows * SERIES_BYTES + query.len() as f64; + detail = json!({"scanned_samples":samples,"syntax_objects":operators,"cpu_seconds_per_item":CPU_PER_ITEM, + "formula":"(samples + 1) * log2(max(samples, 2)) * syntax_objects * cpu_seconds_per_item"}); + (samples + 1.0) * samples.max(2.0).log2() * operators * CPU_PER_ITEM + } + Node::Logical { + operator: Op::CurrentSeries { .. }, + .. + } => { + output_rows = cardinality; + if entry.population_snapshot().is_some() { + let (cpu, workspace) = snapshot_work(entry, cardinality)?; + detail = json!({"input_rows": cardinality, + "physical_program": entry.physical_dag, + "workspace_bytes_bound": workspace, + "formula": "sum of installed operator work; grouped heap count bounded by input rows", + "cpu_seconds_per_item": CPU_PER_ITEM}); + cpu + } else { + cardinality * CPU_PER_ITEM + } + } + Node::RelationalJoin { + inputs, + join_kind, + pred, + left_schema, + right_schema, + .. + } if matches!( + join_kind, + planner_types::pre_asap::JoinKind::Semi + | planner_types::pre_asap::JoinKind::Anti + ) => + { + output_rows = rows[&inputs[0]]; + let indexed = *join_kind == planner_types::pre_asap::JoinKind::Semi + && serde_json::from_value(pred.clone()).is_ok_and(|predicate| { + asap_physical_operators::dag::planner::equijoin_keys( + &predicate, + left_schema, + right_schema, + ) + .is_ok() + }); + detail = json!({"input_rows": input_rows, "output_rows_upper_bound": output_rows, "indexed_equality_semijoin": indexed}); + if indexed { + input_rows.max(1.0) * input_rows.max(2.0).log2() * CPU_PER_ITEM + } else { + rows[&inputs[0]] * rows[&inputs[1]] * CPU_PER_ITEM + } + } + Node::RelationalJoin { .. } => { + output_rows = input_rows.powi(2); + output_rows * CPU_PER_ITEM + } + Node::ExternalExact { .. } + | Node::Logical { + operator: Op::Scan { .. } | Op::Subquery { .. }, + .. + } => { + return Err(invalid("no automatic model for external exact/generic scan/nested subquery operator")); + } + Node::ReduceSum { grouping, .. } => { + if matches!(grouping, crate::query_plan::PhysicalGrouping::Reduce(keys) if keys.is_empty()) + { + output_rows = 1.0; + } + input_rows * CPU_PER_ITEM + } + Node::ExactReadout { .. } + | Node::Binary { .. } + | Node::Relational { .. } + | Node::Logical { .. } => { + input_rows.max(1.0) * input_rows.max(2.0).log2() * CPU_PER_ITEM + } + }; + let ids: BTreeSet<_> = if matches!( + node, + Node::ReadMaterialization { .. } + | Node::SummaryEstimate { .. } + | Node::SummaryMerge { .. } + ) { + used.iter().flat_map(|s| s.ids.iter().cloned()).collect() + } else { + BTreeSet::new() + }; + let workspace_byte_seconds = if matches!(node, Node::Physical { .. }) { + snapshot_work(entry, input_rows)?.1 * cpu * evaluations + } else if entry.population_snapshot().is_some() { + snapshot_work(entry, cardinality)?.1 * cpu * evaluations + } else { + 0.0 + }; + detail["evaluations"] = json!(evaluations); + detail["output_rows_bound"] = json!(output_rows); + insert( + id, + resources( + cpu * evaluations, + workspace_byte_seconds, + network * evaluations, + ids.into_iter().collect(), + detail, + ), + )?; + rows.insert(node_id, output_rows); + profiles.insert(node_id, used); + } + let id = format!("result:{}", entry.query_id); + let evaluations = manifest.components[&id].occurrences_per_horizon; + let bytes = rows[&entry.root] * SERIES_BYTES * evaluations; + insert( + id, + resources( + bytes * CPU_PER_BYTE, + 0.0, + bytes, + vec![], + json!({"rows_per_evaluation":rows[&entry.root],"evaluations":evaluations,"bytes_per_row":SERIES_BYTES}), + ), + )?; + } + if components.len() != manifest.components.len() { + return Err(invalid( + "automatic cost model did not cover every manifest component", + )); + } + if !components + .values() + .map(ComponentResources::weighted_cost) + .sum::() + .is_finite() + { + return Err(invalid("total workload cost overflow")); + } + Ok(AutomaticCostReport { + model_version: MODEL_VERSION.into(), + weights: json!({"unit":"weighted_resource_seconds","cpu_seconds":1.0,"memory_byte_seconds":MEMORY_WEIGHT,"network_bytes":NETWORK_WEIGHT}), + inputs: json!({"data_workload":data,"horizon_seconds":horizon,"ingestion_rate":rate,"source_cadence_seconds":cadence, + "input_cardinality":cardinality,"capability_snapshot_id":env.capability_snapshot_id,"observed_at_unix_ms":now}), + assumptions, + components, + }) +} +fn tree_size(value: &Value) -> usize { + match value { + Value::Object(fields) => 1 + fields.values().map(tree_size).sum::(), + Value::Array(items) => items.iter().map(tree_size).sum(), + _ => 0, + } +} + +#[cfg(test)] +mod tests { + use super::*; + /// A concrete implementation uses all ERP resources even when dearer; + /// a profile for a different implementation cannot replace its estimate. + #[test] + fn physical_state_resources_prefer_applicable_erp() { + let input: super::super::super::compiler::BackendLocalPlanningInput = serde_json::from_str( + include_str!("../../../../docs/examples/asapquery-planning-snapshot.json"), + ) + .unwrap(); + let (mut request, env) = input.into_physical_compilation_request().unwrap(); + let plan = super::super::super::compiler::DeploymentPlanCompiler + .compile_promql(request.clone(), env) + .unwrap(); + let mut m = plan.precompute_plan.materializations[0].clone(); + m.aggregation_type = A::DatasketchesKLL; + m.parameters = std::collections::HashMap::from([("k".into(), json!(200))]); + let baseline = state(&m, &request, 100).unwrap(); + let mut erp = crate::physical::post_asap::cost_model::tests::erp_cost_fixture(); + erp.artifact.records[0].resources = ErpResourceProfile { + memory_bytes: 1e8, + update_cpu_seconds: 0.1, + merge_cpu_seconds: 0.2, + query_cpu_seconds: 0.3, + }; + let expected = erp.artifact.records[0].resources.clone(); + request.erp = Some(erp); + let measured = state(&m, &request, 100).unwrap(); + assert!(!measured.ids.is_empty()); + assert_eq!(measured.profile, expected); + assert!(measured.profile.memory_bytes > baseline.profile.memory_bytes); + request.erp.as_mut().unwrap().artifact.records[0].implementation = "other-runtime".into(); + let fallback = state(&m, &request, 100).unwrap(); + assert!(fallback.ids.is_empty()); + assert_eq!(fallback.profile, baseline.profile); + } +} diff --git a/control_plane/tests/native_rate_topk.rs b/control_plane/tests/native_rate_topk.rs index 0ac493186..eafc353b8 100644 --- a/control_plane/tests/native_rate_topk.rs +++ b/control_plane/tests/native_rate_topk.rs @@ -232,6 +232,57 @@ fn fixed_window_rate_heap_candidates_install_both_physical_graphs() { ); } +// Candidate costs include the persisted maintenance graph, not just its cheap readout. +#[test] +fn fixed_window_heap_costs_include_native_maintenance() { + let mut wire = serde_json::to_value(fixture(true)).unwrap(); + wire["query_workload"]["repeating_queries"][0]["demand"]["fixed_interval_at"]["interval"] = + 60_000.into(); + wire["query_workload"]["repeating_queries"][0]["demand"]["fixed_interval_at"] + ["evaluation_phase"] = 0.into(); + let input: BackendLocalPlanningInput = serde_json::from_value(wire).unwrap(); + let (request, environment) = input.clone().into_physical_compilation_request().unwrap(); + let ids = enumerate_exact_and_materialized_candidates(request) + .unwrap() + .into_iter() + .filter_map(|candidate| { + DeploymentPlanCompiler + .compile_promql(candidate, environment.clone()) + .ok() + }) + .filter(|plan| { + plan.precompute_plan + .executable_dags + .values() + .any(|dag| !dag.native_programs.is_empty()) + }) + .map(|plan| plan.envelope.plan_id) + .collect::>(); + assert!(!ids.is_empty()); + let report = input.compile_promql().unwrap().cost_comparison.unwrap(); + let candidates = report + .candidate_evaluations + .iter() + .filter(|candidate| candidate.plan_id.is_some_and(|id| ids.contains(&id))) + .collect::>(); + assert!(!candidates.is_empty()); + for candidate in candidates { + assert!( + candidate.total_cost.is_some(), + "{:?}", + candidate.unavailable_reason + ); + let resources = candidate.automatic_cost.as_ref().unwrap(); + assert!(resources + .components + .iter() + .any(|(id, resource)| id.starts_with("maintenance:") + && resource.cpu_seconds > 0. + && resource.memory_byte_seconds > 0. + && resource.calculation.get("physical_program").is_some())); + } +} + // Rate must precede grouped Sum in both placements; deployment chooses ownership. #[test] fn grouped_rate_has_native_query_and_maintenance_candidates() { @@ -243,9 +294,10 @@ fn grouped_rate_has_native_query_and_maintenance_candidates() { entry["demand"]["fixed_interval_at"]["evaluation_phase"] = 0.into(); wire["implementation"]["accuracy_evidence"] = json!({}); let input: BackendLocalPlanningInput = serde_json::from_value(wire).unwrap(); - let (request, environment) = input.into_physical_compilation_request().unwrap(); + let (request, environment) = input.clone().into_physical_compilation_request().unwrap(); let mut placements = std::collections::BTreeSet::new(); let mut errors = Vec::new(); + let mut ids = std::collections::BTreeSet::new(); for candidate in enumerate_exact_and_materialized_candidates(request).unwrap() { let plan = match DeploymentPlanCompiler.compile_promql(candidate, environment.clone()) { Ok(plan) => plan, @@ -269,10 +321,23 @@ fn grouped_rate_has_native_query_and_maintenance_candidates() { .any(|dag| !dag.native_programs.is_empty()); assert_eq!(program.to_string().contains("SummaryBuild"), !stored); placements.insert(stored); + ids.insert(plan.envelope.plan_id); } assert_eq!( placements, std::collections::BTreeSet::from([false, true]), "{errors:#?}" ); + let report = input.compile_promql().unwrap().cost_comparison.unwrap(); + for candidate in report + .candidate_evaluations + .iter() + .filter(|candidate| candidate.plan_id.is_some_and(|id| ids.contains(&id))) + { + assert!( + candidate.total_cost.is_some(), + "{:?}", + candidate.unavailable_reason + ); + } } diff --git a/control_plane/tests/native_snapshot_topk.rs b/control_plane/tests/native_snapshot_topk.rs index c4a03c7b1..9eae454b1 100644 --- a/control_plane/tests/native_snapshot_topk.rs +++ b/control_plane/tests/native_snapshot_topk.rs @@ -133,3 +133,48 @@ fn unknown_snapshot_heap_guarantee_cannot_be_installed() { } } } + +// Analytical costing must account for transient heap state and operator work, +// rather than pricing every population program as the same Sort/Limit pair. +#[test] +fn automatic_costs_include_installed_heap_workspace() { + let input = fixture(true); + let (request, environment) = input.clone().into_physical_compilation_request().unwrap(); + let heap_ids = enumerate_exact_and_materialized_candidates(request) + .unwrap() + .into_iter() + .filter_map(|candidate| { + DeploymentPlanCompiler + .compile_promql(candidate, environment.clone()) + .ok() + }) + .filter(heap) + .map(|plan| plan.envelope.plan_id) + .collect::>(); + assert!(!heap_ids.is_empty()); + let plan = input.compile_promql().unwrap(); + let report = plan.cost_comparison.unwrap(); + let evaluated = report + .candidate_evaluations + .iter() + .filter(|candidate| candidate.plan_id.is_some_and(|id| heap_ids.contains(&id))) + .collect::>(); + assert!(!evaluated.is_empty()); + for candidate in evaluated { + assert!( + candidate.total_cost.is_some(), + "{:?}", + candidate.unavailable_reason + ); + let resources = candidate.automatic_cost.as_ref().unwrap(); + assert_eq!(resources.model_version, "backend-workload-resources-v2"); + assert!(resources.components.values().any(|component| { + component.calculation.get("physical_program").is_some() + && component.calculation["workspace_bytes_bound"] + .as_f64() + .unwrap_or(0.) + > 0. + && component.memory_byte_seconds > 0. + })); + } +} diff --git a/data_plane/tests/support/distinct_planning_process.rs b/data_plane/tests/support/distinct_planning_process.rs index 72bb3dd74..242844ec5 100644 --- a/data_plane/tests/support/distinct_planning_process.rs +++ b/data_plane/tests/support/distinct_planning_process.rs @@ -1,4 +1,5 @@ use super::*; +use control_plane::physical::compiler::BackendLocalPlanningInput; /// Modeled HLL error without a confidence certificate cannot replace exact execution. #[tokio::test] @@ -14,3 +15,181 @@ async fn uncertified_distinct_uses_exact_process() { fixture["query_workload"]["repeating_queries"] = serde_json::json!([entry]); assert_uncertified_exact_process(fixture, &[QUERY]).await; } + +/// The production compiler, ingest engine and query DAG preserve distinct populations. +#[tokio::test] +async fn bounded_classic_hll_executes_with_source_labels() { + const QUERY: &str = "distinct_over_time(distinct_values{job=\"api\"}[5s])"; + let fallback_listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let fallback_url = format!("http://{}", fallback_listener.local_addr().unwrap()); + let fallback_task = tokio::spawn(async move { + axum::serve( + fallback_listener, + Router::new().route("/-/healthy", get(|| async { "healthy" })), + ) + .await + .unwrap(); + }); + let mut fixture: Value = serde_json::from_str(include_str!( + "../../../docs/examples/asapquery-compatibility-demo-snapshot.json" + )) + .unwrap(); + let mut entry = fixture["query_workload"]["repeating_queries"][3].clone(); + entry["query"] = QUERY.into(); + entry["requirements"]["accuracy"] = serde_json::json!({"explicit": {"Epsilon": 0.05}}); + fixture["query_workload"]["repeating_queries"] = serde_json::json!([entry]); + // The generated finite value domain enforces this bound over complete + // readout populations, including all panes. Series count is not the bound. + fixture["implementation"]["data_snapshot_id"] = "process-fixture".into(); + fixture["implementation"]["accuracy_evidence"] = serde_json::json!({QUERY: { + "query_string": QUERY, "data_snapshot_id": "process-fixture", + "data_workload": fixture["data_workload"], "source": "finite-generated-value-domain", + "observed_at_unix_ms": fixture["environment"]["observed_at_unix_ms"], + "valid_for_ms": 60000, + "hll": {"model":"asap-classic64-uniform-hash-linear-counting-v1", + "max_distinct_per_readout":128} + }}); + let plan = quote_snapshot_for_test( + serde_json::from_value::(fixture.clone()).unwrap(), + ) + .compile_promql() + .unwrap(); + assert_eq!(plan.precompute_plan.materializations.len(), 1); + assert_eq!( + plan.precompute_plan.materializations[0].aggregation_type, + asap_types::AggregationType::HLL + ); + eprintln!( + "DISTINCT_PLANNED {}", + serde_json::json!({"materializations": plan.precompute_plan.materializations, "query_plan": plan.query_plan, "lifecycle_estimates": plan.lifecycle_estimates}) + ); + let output = tempfile::tempdir().unwrap(); + let path = output.path().join("planning.json"); + let priced = quote_snapshot_for_test(serde_json::from_value(fixture.clone()).unwrap()); + std::fs::write(&path, serde_json::to_vec(&priced).unwrap()).unwrap(); + let port = unused_port(); + let mut vm_port = unused_port(); + while vm_port == port { + vm_port = unused_port(); + } + let mut child = ChildGuard( + Command::new(env!("CARGO_BIN_EXE_data_plane")) + .args([ + "--forward-unsupported-queries", + "--prometheus-server", + &fallback_url, + "--profile", + "asapquery", + "--planning-snapshot", + ]) + .arg(&path) + .args(["--http-port", &port.to_string(), "--output-dir"]) + .arg(output.path()) + .args([ + "--victoriametrics-http-port", + &vm_port.to_string(), + "--victoriametrics-url", + &fallback_url, + ]) + .args([ + "--precompute-allowed-lateness-ms", + "0", + "--precompute-flush-interval-ms", + "25", + ]) + .stdout(Stdio::null()) + .stderr(Stdio::inherit()) + .spawn() + .unwrap(), + ); + let client = reqwest::Client::new(); + let backend = format!("http://127.0.0.1:{port}"); + wait_until_ready(&client, &format!("{backend}/api/v1/health"), &mut child.0).await; + // Source syntax uses the shared parser fork; serving semantics and exact + // routing belong to the MetricsQL adapter and its installed query entries. + let snapshot = serde_json::from_value::(fixture).unwrap(); + let mut snapshot = snapshot; + snapshot.environment.plan_version = 2; + let compiled = quote_snapshot_for_frontend_test(snapshot, true) + .compile_metricsql() + .unwrap(); + let identity = serde_json::json!({"plan_id": compiled.envelope.plan_id, "plan_version": compiled.envelope.plan_version}); + let install = data_plane::drivers::query::servers::http::PhysicalPlanInstallRequest { + summary_catalog: compiled.summary_catalog, + collector_plans: compiled.collector_plans, + precompute_plan: compiled.precompute_plan, + transmission_plan: compiled.transmission_plan, + query_plan: compiled.query_plan, + storage_routing: None, + adaptation_evidence: vec![], + }; + eprintln!( + "DISTINCT_INSTALLED {}", + serde_json::to_string(&install).unwrap() + ); + let response = client + .post(format!("{backend}/api/v1/physical-plan")) + .json(&install) + .send() + .await + .unwrap(); + assert!( + response.status().is_success(), + "{}", + response.text().await.unwrap() + ); + let response = client + .post(format!("{backend}/api/v1/physical-plan/activate")) + .json(&identity) + .send() + .await + .unwrap(); + assert!( + response.status().is_success(), + "{}", + response.text().await.unwrap() + ); + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_millis() as i64; + let base = now - now.rem_euclid(5000) - 20000; + let mut series = Vec::new(); + for (instance, distinct, job) in [("a", 5, "api"), ("b", 13, "api"), ("excluded", 23, "other")] + { + let mut samples: Vec<_> = (0..100) + .map(|i| (base + 1 + i, (i % distinct) as f64)) + .collect(); + samples.push((base + 15001, 1000.0)); + series.push(series_with_labels( + "distinct_values", + &[("instance", instance), ("job", job)], + &samples, + )); + } + assert_eq!( + remote_write(&client, &backend, &WriteRequest { timeseries: series }).await, + 204 + ); + drain_precompute(&client, &backend).await; + let result = wait_for_warm_instant( + &client, + &format!("http://127.0.0.1:{vm_port}"), + QUERY, + (base + 5000) as f64 / 1000.0, + &output.path().join("query_engine.log"), + ) + .await; + let rows = result["data"]["result"].as_array().unwrap(); + assert_eq!(rows.len(), 2, "{result}"); + for (instance, exact) in [("a", 5.0), ("b", 13.0)] { + let row = rows + .iter() + .find(|row| row["metric"]["instance"] == instance) + .unwrap(); + let estimate = row["value"][1].as_str().unwrap().parse::().unwrap(); + assert!((estimate - exact).abs() / exact <= 0.05, "{result}"); + } + eprintln!("DISTINCT_WARM {result}"); + fallback_task.abort(); +} diff --git a/docs/design_docs/asapplanner-integration.md b/docs/design_docs/asapplanner-integration.md index 6d5332ad1..531f2e3b0 100644 --- a/docs/design_docs/asapplanner-integration.md +++ b/docs/design_docs/asapplanner-integration.md @@ -82,6 +82,13 @@ A maintenance lifecycle is a contract associated with computation, not another operator IR. A deployment plan is an operational wrapper around Physical DAGs, not another lowering stage. +Cost evidence includes initialization, ingestion updates, retained state, +transmission, storage, merges, readouts and recurring queries, counting shared +producer construction once. Candidates use the same data and demand scope; +missing evidence is not zero cost. Backend `candidate_cost()` supplies applicable +ERP or analytical estimates to Planner selection. See the +[workload calculation model](evidence-dependent-candidates.md#backend-owned-workload-calculation). + ## 3. Deployment plan structure One installed version contains: diff --git a/docs/design_docs/evidence-dependent-candidates.md b/docs/design_docs/evidence-dependent-candidates.md index 9e09b1f86..d4baf36cc 100644 --- a/docs/design_docs/evidence-dependent-candidates.md +++ b/docs/design_docs/evidence-dependent-candidates.md @@ -1,7 +1,6 @@ # Evidence-dependent candidate selection and deployment -Status: scoped evidence and native candidate admission implemented by backend -PR #761, building +Status: evidence and workload costing implemented by backend PR #761, building on the API adaptation in #768, against ASAPPlanner #455 (`2ec3fc80`). This document defines the backend decision boundary for Planner issue #454 and backend issue #752. It does not claim that every retained @@ -19,7 +18,7 @@ could prove valid or deploys one whose guarantee has never been established. | Construct semantic candidates | Planner | Preserve unknown guarantees; reject known-invalid evidence and impossible shapes | | Supply external facts | Backend | Bind evidence to the query, data population, snapshot and validity period | | Derive accuracy and select logical roots | Planner, under backend models and policy | Respect the root accuracy target; missing proof is not certification | -| Bind and admit a deployment | Backend | Verify concrete execution support and require complete workload cost evidence | +| Bind and admit a deployment | Backend | Verify concrete execution support and compute complete workload resource cost | The backend still invokes Planner's workload search and global selection. It does not introduce a second semantic optimizer. Physical binding preserves the @@ -38,7 +37,7 @@ flowchart TD Search --> Select[Planner global selection under backend policy] Select --> Exact[Explicit exact fallback when no certified summary is selected] Select --> Bind[Bind selected logical DAG to concrete execution] - Bind --> Admit[Check guarantees, runtime support and complete workload cost] + Bind --> Admit[Check guarantees, runtime support and computed workload cost] Admit -->|Pass| Publish[Publish coherent physical plan] Admit -->|Fail| Reject[Reject deployment] ``` @@ -75,6 +74,31 @@ The backend validates the supplied record's scope and consistency. It does not derive a source-domain proof from samples or establish the truth of a producer's claimed contract. Producing valid external proofs remains an upstream duty. +### Bounded classic HLL confidence + +Scoped evidence may carry `hll: {"model": +"asap-classic64-uniform-hash-linear-counting-v1", "max_distinct_per_readout": 128}`. +This declares an enforced source-domain upper bound for **each complete readout +population**, including the union of all merged panes, and the independent +uniform bucket-hash assumption. Series count, a sampled distinct count, and +ERP observed maxima do not establish this contract. The bound is supplied by +the source-contract owner; the backend validates its query/data/snapshot/time +scope, not the truth of the external source assertion. + +The model supports declared upper bounds from 1 through 4096 and precisions +4 through 18. It only certifies parameters that keep every permitted population +in classic HLL's linear-counting branch. A finite collision-arrival bound gives +failure probability for the requested relative error. Sizing searches for the +smallest supported precision meeting the same epsilon/delta contract. Explain +retains the model, population limit, hash assumption and numerical guarantee. + +The implementation is bound to the backend-local Regular HLL estimator; +collector estimators, HIP, MLE and unbounded populations are not certified by it. +Missing contracts retain generic HLL's unknown probability; infeasible targets +retain exact execution. ERP resources can still price the selected parameters, +but ERP error maxima cannot override this confidence model. Each readout's +probability is not a simultaneous guarantee for an entire dashboard. + ## Accuracy, runtime and cost remain separate A known accuracy guarantee does not establish executor support. The physical @@ -84,8 +108,8 @@ must not be reported as deployment approval. Missing numerical cost remains unavailable, never zero. The backend implements Planner's `candidate_cost()` directly; it does not enable uncosted legacy -selection. A cost estimate only supports logical selection. Publication requires -complete, applicable workload cost quotes. A cheap candidate cannot bypass +selection. A cost estimate only supports logical selection. Publication requires complete backend-computed resource costs or an explicit, +applicable provider override. A cheap candidate cannot bypass accuracy or runtime admission. Explain records accuracy status and symbolic guarantee, runtime support status, @@ -94,6 +118,168 @@ is labelled as pending backend binding. Rejected candidates retain their reported reasons. This distinguishes missing proof, unsupported execution, missing comparable cost and a candidate that simply lost the ranking. +## ERP and analytical cost models + +ERP is the benchmark source for this path. The existing ERP artifact and +runtime-observation input feed both parameter planning and resource estimation; +there is no separate benchmark upload contract for candidate costs. + +```mermaid +flowchart TD + Profile[ERP benchmark artifact] --> Match[Match implementation, exact parameters and data population] + Observations[Declared distribution or validated runtime shape] --> Match + Candidate[Planner candidate with concrete parameters] --> Match + Match -->|Applicable profile| Measured[ERP resource estimate] + Match -->|No applicable profile| Analytical[Backend analytical resource estimate] + Measured --> Cost[Backend candidate_cost and source explanation] + Analytical --> Cost + Cost --> Selection[Planner logical selection] + Proof[Accuracy evidence and guarantee model] --> Selection + Selection --> Physical[Physical binding and complete workload costing] + Physical --> Admission[Deployment admission] +``` + +**Priority is applicable ERP, then analytical, then unavailable.** A larger +measured value still overrides a smaller analytical estimate. ERP matching +uses its existing implementation/runtime filters, exact sketch parameters, +minimum trial count, distribution equality or configured bounded shape match. +Catalog-scoped observations retain their existing population and freshness +validation. A profile for another parameter point or implementation cannot be +substituted. ERP v1 artifacts themselves have no per-record expiry field; do +not confuse runtime-observation freshness with a benchmark expiry guarantee. + +Cost lookup does not certify empirical accuracy: measured resource usage can +be useful even when observed error does not meet a target or cannot establish +its failure probability. The cost path reuses ERP profile matching without an +empirical-error threshold; the accuracy path separately checks the requested +guarantee. In particular, an explicit confidence target must not erase ERP +resource measurements merely because ERP v1 cannot prove that confidence. +Existing empirical-only accuracy policy still applies to accuracy decisions. + +The first backend model, `backend_state_footprint_v1`, estimates **retained +state bytes per partition** for local candidate selection. It uses ERP's +`memory_bytes` when applicable; otherwise it reuses the backend's existing +analytical retained-state formulas: matrix dimensions and heap capacity, +KLL capacity, HLL registers, fixed accumulator allowance and the current +DDSketch allowance. These are estimates, not measured limits or a complete +workload cost. Reachable shared state nodes are counted once. For multiple +observed populations the ERP proxy uses the largest matched partition, without +pooling the populations. Shared-grid families need their own applicable model; +an independent sketch profile is not a Hydra-grid measurement. + +The estimate excludes population counts, pane multiplicity, CPU, transmission +and other deployment costs. An applicable ERP profile also ranks sketch-family +candidates using the same byte estimates, with analytical estimates for the +unmeasured peers. Without an applicable profile, existing family preference +order remains the fallback policy. +It does not mix ERP CPU seconds with analytical bytes or claim to minimize +complete workload cost. Raw rewrites and exact compositions have no state-byte +estimate here; exact composition retains its separate measured recurring-cost +model. Unsupported shapes remain uncosted. + +Explain attaches the model, unit, source (`erp`, `analytical`, or `mixed`) and +ERP record IDs to each available estimate. Final deployment selection combines these unit resources with physical demand +as described below. A profile alone neither prices the workload nor approves +publication. + +Acceptance covers ERP precedence even when measured cost is higher, matching +parameters/implementation/population, insufficient trials, invalid measurements, +analytical fallback, unavailable shapes, explain provenance, and independent +accuracy and publication gates. + +## Backend-owned workload calculation + +Both backend-local startup and the PromQL/MetricsQL compile-and-publish API +compute workload costs when `workload_cost_evidence` is absent. Callers supply +queries, recurrence, data facts, deployment capabilities and optional ERP; they +do not need to manufacture complete candidate quotes. The backend binds the +bounded candidate inventory, constructs its coverage manifest, calculates every +component, and selects the lowest comparable total among feasible candidates. +The logical state-footprint proxy above remains a separate early selection +model; physical selection compares only the enumerated candidates, not every +possible logical DAG or window implementation. + +```mermaid +flowchart LR + ERP[Applicable ERP unit resources] --> Units[ERP first; analytical fallback] + Data[Fresh rate, cardinality, cadence, volume] --> Demand[Physical workload quantities] + Queries[Query frequency and common horizon] --> Demand + Plan[Bound DAG, shared state, panes, locations, transport] --> Demand + Units --> Compute[CPU seconds, memory byte-seconds, network bytes] + Demand --> Compute + Compute --> Compare[Common weights and complete component coverage] + Compare --> Select[Lowest-cost feasible physical candidate] +``` + +`backend-workload-resources-v2` is an explicit analytical reference model, not a +calibrated prediction or a monetary quote. It computes a resource vector before +weighting it: `cost = CPU_seconds + 1e-9 * memory_byte_seconds + 1e-8 * network_bytes`. +These fixed versioned weights express a default tradeoff; a deployment can still +supply complete calibrated provider quotes as an explicit override. Override +quotes are validated as one model across the inventory; missing/invalid quotes +are not silently patched with analytical values in unrelated units. + +For horizon H, input rate R, query interval I seconds and partition estimate P: + +| Component | Backend calculation | +| --- | --- | +| Source | R × H samples per distinct source/location; exact sources also retain raw samples for the required lookback plus series metadata | +| State initialization/retirement | P × (initial retained states + state rotations during H); allocation charges estimated bytes, retirement charges state count | +| State updates | R × H × overlapping full windows; panes receive each input once; rollup levels add their child-to-parent merges | +| State residency | P × retained physical states × unit state bytes × H | +| Transport | Rule emission/checkpoint count × partitions × retained states × estimated encoded bytes, plus serialization, decoding and receiving merges | +| Summary reads | H / I × P × panes read; decode each pane, merge additional panes, then charge the selected state's query/readout unit cost | +| Other query operators | Traverse each reachable physical node once per evaluation, using propagated row estimates and an explicit analytical operation model | +| Native exact subtrees | Source sample count for the lookback or full data volume, syntax complexity and a sorting allowance; charge each subtree's returned data separately from final output | +| Result | Result row estimate × row bytes × H / I, including serialization and delivery | + +ERP supplies state memory and per-update, per-merge and per-query CPU seconds. +The backend uses the actual workload quantities, not ERP's example invocation +counts. Applicability checks are the same implementation/parameter/runtime/ +population/shape/trial checks used by candidate costs. Multiple matched +population profiles use the maximum of each resource dimension per partition. +ERP measurements override analytical estimates even when more expensive. They +do not establish accuracy guarantees or runtime availability. + +Unmeasured dimensions use versioned analytical assumptions: 1e-7 CPU seconds per +item, 1e-9 CPU seconds per byte, 24 bytes per raw sample, and a 256-byte series/ +result-row allowance. Sketch update work scales with log2(state bytes / 16), +while merge/read work scales with state bytes. Native exact work uses +`(samples + 1) × log2(max(samples, 2)) × canonical syntax object count × item CPU`. +This is a transparent complexity proxy, not a benchmark of Prometheus; local +arithmetic/reduction/sort and joins have their own row-based estimates. + +The current data contract has workload-level facts, not a complete per-source +histogram. The model explicitly replicates the workload rate/cardinality to +each distinct source and collector rather than claiming known selectivities. +Grouped states use input cardinality as a group-count estimate. Without a fresh +cardinality, it estimates `ceil(R × scrape_interval)` under the stated assumption +of one sample per active series per scrape; a zero rate cannot establish an +unknown population. At-rest/full-history costing requires fresh ingestion volume. +Unknown rate/cadence, unsupported operators, invalid values and overflow make a +candidate unavailable. No partial total is admitted. + +Shared state and source upkeep are charged once per physical identity/location; +additional consumers add reads and outputs. For distributed plans, receiving +backend merges are charged in transport rather than again as raw updates. +Full-window overlap and pane retention are taken from the compiled layout; +rollup reads conservatively use base panes. Delta transport conservatively uses +full payload size and checkpoints; there is no unmeasured compression discount. +All recurring work uses the same horizon and original query demand. Existing +window implementation/lifecycle inputs still determine which physical layouts +are bound; their opaque weighted costs are not mixed with this resource model. + +Each candidate's `automatic_cost` report records the input data, decision time, +model/weights, assumptions, component resource vectors, workload multipliers and +applicable ERP record IDs. `component_costs` is the weighted projection of that +report, and covers exactly the existing manifest. Invalid provider overrides +remain explicit errors. Accuracy proofs, compiler capability checks and runtime +readiness/fallback policy remain independent admission requirements. + +Acceptance includes quote-free startup and HTTP compilation, ERP precedence and +implementation mismatch, frequency/data-size scaling, shared-state deduplication, +expired facts, overflow, complete coverage and provider-override compatibility. + ## Example and acceptance behavior For `quantile_over_time(0.9, data[5m]) / quantile_over_time(0.5, data[5m])`: @@ -103,8 +289,8 @@ For `quantile_over_time(0.9, data[5m]) / quantile_over_time(0.5, data[5m])`: 2. With valid, scoped operand contracts, Planner can derive the ratio guarantee and check the explicit root target. Valid evidence alone does not guarantee that the target is met or that this candidate wins selection. -3. A selected candidate still needs physical support and complete workload cost - quotes before publication. A stale or cross-query certificate rejects the request. +3. A selected candidate still needs physical support and complete backend-computed + workload costs before publication. A stale or cross-query certificate rejects the request. Cross-family acceptance includes absent, partial, valid, invalid and stale evidence, an explicit root accuracy target, unavailable cost, and unsupported @@ -136,14 +322,16 @@ Spatial TopK admission keeps the exact population ranking and the Planner's CountSketch-with-heap physical candidate separate. The heap requires scoped score separation and an enforced `topk_max_distinct_items` bound; an estimated workload cardinality is insufficient. Missing evidence leaves the candidate -visible but ineligible for installation. A complete provider quote can +visible but ineligible for installation. Cost model v2 prices the installed +native operators and transient heap workspace. A complete provider quote can choose either candidate; neither is forced by operator name. Rate TopK follows the same admission and pricing boundary. Nonnegative finalized counter rates admit CMS as well as CountSketch when the scoped accuracy evidence is sufficient. Exact ranking remains available. The installed physical program consumes the complete per-series Rate vector from bound counter SDS; Backend -quotes can select any of the three programs. Operator support alone does not supply accuracy proof. +quotes can select any of the three programs, and automatic costs charge the +actual native workspace. Operator support alone does not supply accuracy proof. A fixed-window Rate heap candidate places Rate finalization and heap construction in precompute and persists the heap. Query-time construction remains a separate diff --git a/docs/design_docs/shape-aware-erp-v1.md b/docs/design_docs/shape-aware-erp-v1.md index 5d04d99af..0f0a7fccf 100644 --- a/docs/design_docs/shape-aware-erp-v1.md +++ b/docs/design_docs/shape-aware-erp-v1.md @@ -22,8 +22,12 @@ than publishing a biased partial snapshot. Profiles with too few benchmark events, poor fit, ambiguous confidence, or excessive cardinality/parameter distance are misses. -On a hit, empirical parameters and measured atomic costs are used. On a miss, -malformed evidence, or drift, Hybrid mode retains the theoretical parameters; +On a hit, empirical parameters and measured atomic costs inform planning. +ERP resource precedence, analytical fallback, accuracy certification and +deployment admission follow the +[evidence-dependent candidate design](evidence-dependent-candidates.md#erp-and-analytical-cost-models). + +On a miss, malformed evidence, or drift, Hybrid mode retains the theoretical parameters; if the runtime cannot deploy them or they exceed its memory limit, compilation chooses exact execution. Empirical-only mode fails closed. diff --git a/docs/developer_docs/query-engine/catalog-physical-plan-runtime.md b/docs/developer_docs/query-engine/catalog-physical-plan-runtime.md index cf575abbd..885b56bfb 100644 --- a/docs/developer_docs/query-engine/catalog-physical-plan-runtime.md +++ b/docs/developer_docs/query-engine/catalog-physical-plan-runtime.md @@ -94,7 +94,8 @@ For selected explicit-`by` current-series TopK, the population store now supplie all eligible members through a snapshot binding. Planner compiles the ranking above that boundary; the QueryPlan persists that physical program before candidate pricing and activation. Serving recovers its operators and supplies -full-label native rows without parsing or lowering the query. Provider quote +full-label native rows without parsing or lowering the query. Automatic cost +estimates include native operator CPU and temporary workspace; provider quote manifests include the physical program itself. Other population readouts retain their existing paths. diff --git a/docs/evaluation/e2e-physical-dag.md b/docs/evaluation/e2e-physical-dag.md index 4a91a4aeb..d781e7061 100644 --- a/docs/evaluation/e2e-physical-dag.md +++ b/docs/evaluation/e2e-physical-dag.md @@ -39,9 +39,10 @@ Planner owns query semantics, summary families, parameters and candidate selection. ERP can affect supported evidence-based choices, but supplying an artifact does not prove that it was eligible or used. Freshness, source/update semantics, parameters and accuracy constraints still apply. The checked-in demo -is not an empirical-ERP benchmark. For deployment, use the priced snapshot workflow in [execution calibration](../../tools/o11y-execution/CALIBRATION.md). -`compile_workload_artifact` requires complete workload quotes; do not relabel -demo costs as measured evidence. +is not an empirical-ERP benchmark. For deployment, the backend automatically computes workload costs from ERP or +analytical unit resources and physical demand. Optional calibrated overrides use +the [cost evidence workflow](../examples/workload-cost-evidence.md). Do not relabel +analytical demo costs as measured evidence. For the existing observation → ERP-selected KLL → installed HTTP correctness fixture, see [ERP process validation](../developer_docs/erp-process-validation.md): @@ -73,12 +74,11 @@ dot -Tsvg target/physical-dag-inspection/selected.dot \ -o target/physical-dag-inspection/selected.svg ``` -`compile_workload_artifact` calls the same evidence-required snapshot compiler -used by startup. Set `ASAPQUERY_PLANNING_SNAPSHOT` to a priced snapshot prepared -using the [cost evidence workflow](../examples/workload-cost-evidence.md). -The checked-in unquoted templates support candidate discovery only. The output -includes the selected plan, logical selection trace and complete cost comparison. -MetricsQL compilation requires quotes collected for that frontend. +`compile_workload_artifact` calls the same snapshot compiler used by startup. +Set `ASAPQUERY_PLANNING_SNAPSHOT` to the workload snapshot. Without an explicit +provider override it uses automatic ERP/analytical workload costing. The output +includes the selected plan, logical selection trace, resource assumptions and +complete cost comparison. Both PromQL and MetricsQL use this path. The compiler derives sibling plans from the selected post-ASAP DAG: diff --git a/docs/examples/workload-cost-evidence.md b/docs/examples/workload-cost-evidence.md index fc5430d55..33e06e5d8 100644 --- a/docs/examples/workload-cost-evidence.md +++ b/docs/examples/workload-cost-evidence.md @@ -1,10 +1,16 @@ # Complete workload cost evidence Planning snapshots use one schema, `snapshot_version: 3`. Versions 1 and 2 are rejected; the environment must declare its logical dataset identity. -Candidate discovery may omit `workload_cost_evidence`; compiling a deployable -snapshot requires complete, valid quotes and selects by complete workload cost. -There is no unquoted snapshot deployment path. The checked-in JSON examples are -discovery templates, not ready-to-deploy plans. +Deployment automatically calculates workload costs when `workload_cost_evidence` +is absent, using applicable ERP profiles and analytical estimates together with +the workload's data facts, query frequency and physical layout. Inspect +`cost_comparison` for candidate totals, resource breakdowns and assumptions. +The checked-in examples can exercise this path at their recorded decision time; +use current data/capability inputs for a live deployment. + +The following workflow is an **optional provider override** for deployments with +calibrated complete quotes. See the [automatic calculation design](../design_docs/evidence-dependent-candidates.md#backend-owned-workload-calculation) +for the default path. ## Workflow @@ -65,7 +71,9 @@ can change which bound workload is committed, not rewrite its semantics. All quotes use one provider model's common cost units. A horizon quote includes the complete stated partition's work, data volume/cardinality, maintained -groups and retention. Do not reuse a global ingestion rate as a per-metric rate. +groups and retention. A calibrated provider should supply per-source demand rather than infer it from +a global rate. The automatic reference model discloses conservative replication +when only workload-level facts are available. Per-evaluation quotes exclude upkeep already charged in horizon components. Sunk infrastructure may explicitly cost zero under the provider's documented decision boundary; unknown costs may not.