diff --git a/crates/devtools/Cargo.toml b/crates/devtools/Cargo.toml index edf20aef..944858a3 100644 --- a/crates/devtools/Cargo.toml +++ b/crates/devtools/Cargo.toml @@ -16,7 +16,7 @@ asap-logical-optimizer = { path = "../logical-optimizer" } # plans with. A devtool, not a stage crate, so the #572 guards allow it. asap-executor = { path = "../executor" } -# Used by the show_ir / dag_export / variant_coverage bins (catalog schemas, +# Used by the show_*_ir / stage_pipeline / variant_coverage bins (catalog schemas, # async SQL path, JSON output) and by the topk_ir / canonical_examples # examples. Regular deps, not dev-deps: `[[bin]]` targets can't see dev-deps. asap-types = { path = "../types" } diff --git a/crates/devtools/src/bin/dag_export.rs b/crates/devtools/src/bin/dag_export.rs deleted file mode 100644 index f94a4078..00000000 --- a/crates/devtools/src/bin/dag_export.rs +++ /dev/null @@ -1,3188 +0,0 @@ -// cargo run -p asap-lower --bin dag_export -- \ -// --sql "SELECT service, COUNT(*) FROM metrics GROUP BY service" --name q1 \ -// --promql "topk(5, rate(http_requests_total[5m]))" --name q2 -// -// Lowers each given SQL/PromQL query to pre-ASAP IR and prints a single -// `asap_types::dag_export::WorkloadDAG` as JSON on stdout — the input format -// for `tools/dag-viewer` (issue #133). Redirect to a file and load it there: -// cargo run -p asap-lower --bin dag_export -- --sql "..." --name q1 > /tmp/dag.json -// -// `--name` is optional; an unnamed query defaults to `q` (1-indexed). -// -// `--epsilon ` is optional and applies to every query in the run: it -// lowers with `AccuracyTarget::Epsilon()` instead of the default -// `AccuracyTarget::Exact`. Without it, every `AggIntent` lowers exact and -// `asap_logical_optimizer::ASAPStrategies` never has a genuine sketch -// alternative to report — so no node ever picks up a `SketchApproximation` -// note. Pass it to actually exercise that path, e.g.: -// cargo run -p asap-lower --bin dag_export -- \ -// --epsilon 0.01 --sql "SELECT quantile(0.99, latency) FROM metrics" --name p99 -// -// `--post-asap` needs one of two cost sources to rank with, and does -// nothing without either (see `--default-cost` and `--planner-cost-json` -// below). -// -// `--post-asap` is optional and off by default. When passed, this binary -// additionally runs `asap_logical_optimizer::pass1::replacement::search_workload` (this -// binary took no strategies of its own — `default_strategies()` already -// includes `AvgToSumOverCountStrategy` as of #282) over every lowered query -// and ranks each discovered `TargetSubDAGCandidates` via `candidate_selection::cost_sorted`. The -// best-ranked -// candidate per group feeds two additive outputs: -// -// - one `asap_types::dag_export::TargetReplacement` per group on whichever -// query's `NamedDAG.replacements` contains that target node (matched -// by `DAGNode::hash` + structural equality, the same collision-safe -// pattern `annotate_with_explanations` below already uses for notes) — -// a small, self-contained "before -> after" pair per replacement site; -// - one merged `NamedDAG.post_dag`: a single flattened DAG per -// query with every winning candidate spliced directly into the query's -// own pre-ASAP shape in place, built via -// `asap_types::dag_export::export_post_asap`. -// -// Together these surface every one of the four concrete replacement kinds: -// the sketch family `ASAPStrategies`/`HydraGroupingStrategy` bound, -// the CSE share/recompute choice `SharedSubDAGStrategy` found, the -// workload-aware roll-up `RollupStrategy` derived, and the `avg -> -// sum/count` rewrite `AvgToSumOverCountStrategy` proposes. Without -// `--post-asap`, every existing invocation of this binary produces -// byte-identical output to before (`NamedDAG.replacements` is empty and -// `post_dag` is `None`, both skipped from the JSON entirely in that -// case). E.g.: -// cargo run -p asap-lower --bin dag_export -- \ -// --post-asap --default-cost --epsilon 0.01 \ -// --sql "SELECT quantile(0.95, latency) FROM metrics" --name q1 -// -// `--default-cost` and `--planner-cost-json` are the two mutually exclusive -// ways to give `--post-asap` a cost model, and they differ in what the -// export is allowed to claim: -// -// - `--planner-cost-json ` supplies complete deployment-owned -// physical-plan evidence. It both ranks the candidates and is exported: -// every decision carries a calibrated `CostUnits` annotation. -// - `--default-cost` ranks with `asap_plan_selection::cost::cost_model:: -// DefaultCostModel` — structural node counts, owning no deployment -// evidence. The structure of the result is real (which replacements the -// search found, which one won per group, what the merged post-ASAP DAG -// looks like); the numbers are not exported at all. Every decision's -// cost is `CostSource::Unavailable` with `value: None`, which the viewer -// renders as "Not estimated". Use it to see what ASAPPlanner does with a -// workload before there is a deployment to measure. -// -// Absent both, `--post-asap` exports the raw plan only. - -use std::cell::RefCell; -use std::collections::HashMap; -use std::rc::Rc; -use std::time::Instant; - -use asap_logical_optimizer::pass1::replacement::{ - default_strategies_with_evidence, is_logical_rewrite, search_workload, search_workload_with, - Replacement, ReplacementSubDAG, -}; -use asap_logical_optimizer::{AccuracyEvidenceProvider, PropagationStats}; -use asap_plan_selection::candidate_selection::global_selection; -use asap_plan_selection::cost::analytical_cost::{ - cache_hit_ratios, AnalyticalCostError, EvidenceBackedPhysicalDAG as PhysicalDAG, - PhysicalNodeEvidence, ResourceCalibration, ANALYTICAL_COST_MODEL_VERSION, -}; -use asap_plan_selection::cost::cost_model::DefaultCostModel; -use asap_plan_selection::cost::cost_model::{Cost, CostModel}; -use asap_plan_selection::cost::physical_operator_statistics::ComparisonScope; -use asap_plan_selection::cost::physical_plan_cost_model::{ - PhysicalEvidenceSnapshot, PhysicalPlanCostModel, PlannerPhysicalPlanProvider, -}; -use asap_plan_selection::cost::query_physical_lowering::PhysicalNodeRequest; -use asap_types::cost::{BaselineRef, CostAnnotation, CostInput, CostSource, CostUnit}; -use asap_types::dag_export::{ - self, DAGDecision, DAGNote, ExportDAG, NamedDAG, PostAsapSubstitution, TargetRejection, - TargetReplacement, TargetReplacementAfter, WorkloadDAG, -}; -use asap_types::ir::cse::{structural_hash, HashCache}; -use asap_types::ir::properties::CompositionOperator; -use asap_types::ir::schema::{DataType, Field, Schema}; -use asap_types::ir::schema::{FieldDataType, SketchStatistic}; -use asap_types::ir::OperatorNode; -use asap_types::types::AccuracyTarget; -use asap_types::workload::resources::CacheProfile; - -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -#[serde(deny_unknown_fields)] -struct PlannerCostDocument { - #[serde(default, skip_serializing_if = "Option::is_none")] - storage_io: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - #[serde(rename = "boundaries")] - handoffs: Option, - /// Immutable catalog/runtime evidence generation shared by this file. - evidence_version: String, - calibration: ResourceCalibration, - targets: Vec, -} - -fn parse_planner_cost_document(raw: &str) -> Result { - let value: serde_json::Value = serde_json::from_str(raw) - .map_err(|error| format!("planner cost evidence is invalid JSON: {error}"))?; - if value.get("targets").is_none() { - return Err("the compact analytical-cost payload is no longer supported; migrate to complete per-operator physical DAG evidence".into()); - } - let document: PlannerCostDocument = serde_json::from_value(value) - .map_err(|error| format!("planner cost evidence is invalid: {error}"))?; - if document.evidence_version.trim().is_empty() { - return Err("planner cost evidence_version must be non-empty".into()); - } - Ok(document) -} - -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -#[serde(deny_unknown_fields)] -struct TargetPhysicalEvidence { - target: Rc, - scope: ComparisonScopeEvidence, - candidates: Vec, -} - -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -#[serde(deny_unknown_fields)] -struct ComparisonScopeEvidence { - data_arrival: asap_types::workload::DataArrival, - planning_time_ms: u64, - horizon_ms: u64, - evaluation_count: u64, - time_scope: String, - lookback_ms: Option, - as_of_ms: Option, - sources: Vec, - #[serde(default = "CacheProfile::no_cache")] - cache_profile: CacheProfile, -} - -impl ComparisonScopeEvidence { - fn resolve(&self) -> Result { - use asap_types::workload::{ - DurationMs, QueryRecurrence, QueryTimeScope, TimeSelection, TimestampMs, - }; - let scope = match self.time_scope.as_str() { - "real_time" => QueryTimeScope::RealTime, - "longitudinal" => QueryTimeScope::Longitudinal, - "mixed" => QueryTimeScope::Mixed, - "unknown" => QueryTimeScope::Unknown, - _ => return Err(AnalyticalCostError::MissingComparisonScope("time_scope")), - }; - let comparison = ComparisonScope { - data_arrival: self.data_arrival, - planning_time: TimestampMs(self.planning_time_ms), - horizon: DurationMs(self.horizon_ms), - recurrence: QueryRecurrence::OneTime { - invocations: self.evaluation_count, - execute_at: None, - }, - time_selection: TimeSelection { - scope, - lookback: self.lookback_ms.map(DurationMs), - as_of: self.as_of_ms.map(TimestampMs), - }, - sources: self.sources.clone(), - }; - comparison.validate()?; - Ok(comparison) - } -} - -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -#[serde(deny_unknown_fields)] -struct QueryNodePhysicalEvidence { - logical_node: OperatorNode, - operator: asap_plan_selection::cost::analytical_cost::PhysicalOperator, - occurrence: usize, - synthetic: bool, - evidence: PhysicalNodeEvidence, -} - -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)] -enum CandidatePhysicalEvidence { - Rewrite { - plan: serde_json::Value, - query_nodes: Vec, - }, - Summary { - plan: serde_json::Value, - query_nodes: Vec, - physical_dag: PhysicalDAG, - }, -} - -impl CandidatePhysicalEvidence { - fn plan(&self) -> &serde_json::Value { - match self { - Self::Rewrite { plan, .. } | Self::Summary { plan, .. } => plan, - } - } - - fn query_nodes(&self) -> &[QueryNodePhysicalEvidence] { - match self { - Self::Rewrite { query_nodes, .. } | Self::Summary { query_nodes, .. } => query_nodes, - } - } - - fn matches(&self, candidate: &ReplacementSubDAG) -> bool { - let actual = match (self, &candidate.replacement) { - (Self::Summary { .. }, Replacement::SubDAG(node)) if !is_logical_rewrite(node) => { - serde_json::to_value(dag_export::export(node)) - } - (Self::Rewrite { .. }, Replacement::SubDAG(node)) if is_logical_rewrite(node) => { - serde_json::to_value(dag_export::export(node)) - } - _ => return false, - }; - actual.is_ok_and(|actual| plan_values_match(&actual, self.plan())) - } - - fn summary_dag(&self) -> Option<&PhysicalDAG> { - match self { - Self::Summary { physical_dag, .. } => Some(physical_dag), - Self::Rewrite { .. } => None, - } - } -} - -/// Compare complete exported-plan identity while tolerating the one-ULP -/// decimal round trip that `serde_json::Value` can introduce for derived -/// floating-point guarantees. Integer configuration and DAG identity stay -/// exact; no guarantee field is dropped or otherwise normalized away. -fn plan_values_match(actual: &serde_json::Value, expected: &serde_json::Value) -> bool { - plan_values_match_inner(actual, expected, false) -} - -fn plan_values_match_inner( - actual: &serde_json::Value, - expected: &serde_json::Value, - inside_guarantee: bool, -) -> bool { - use serde_json::Value; - match (actual, expected) { - (Value::Object(actual), Value::Object(expected)) => { - actual.len() == expected.len() - && actual.iter().all(|(key, value)| { - expected.get(key).is_some_and(|other| { - plan_values_match_inner( - value, - other, - inside_guarantee || key == "guarantee", - ) - }) - }) - } - (Value::Array(actual), Value::Array(expected)) => { - actual.len() == expected.len() - && actual - .iter() - .zip(expected) - .all(|(value, other)| plan_values_match_inner(value, other, inside_guarantee)) - } - (Value::Number(actual), Value::Number(expected)) => { - match (actual.as_i64(), expected.as_i64()) { - (Some(actual), Some(expected)) => actual == expected, - _ => match (actual.as_u64(), expected.as_u64()) { - (Some(actual), Some(expected)) => actual == expected, - _ => match (actual.as_f64(), expected.as_f64()) { - (Some(actual), Some(expected)) => { - actual == expected - || (inside_guarantee - && actual.to_bits().abs_diff(expected.to_bits()) <= 1) - } - _ => false, - }, - }, - } - } - _ => actual == expected, - } -} - -struct ExportPhysicalProvider<'a> { - storage_io: Option<&'a asap_plan_selection::cost::storage_io::StorageIoProfile>, - handoffs: Option<&'a asap_plan_selection::cost::physical_handoff_cost::PhysicalHandoffProfile>, - evidence_version: &'a str, - target: &'a TargetPhysicalEvidence, - candidate: &'a CandidatePhysicalEvidence, - used_query_nodes: RefCell>, -} - -impl ExportPhysicalProvider<'_> { - fn all_query_evidence_used(&self) -> bool { - self.used_query_nodes.borrow().len() == self.candidate.query_nodes().len() - } -} - -impl PlannerPhysicalPlanProvider for ExportPhysicalProvider<'_> { - fn capture_evidence_snapshot( - &self, - _target: &asap_logical_optimizer::pass1::replacement::TargetSubDAG<'_>, - ) -> Result { - Ok(PhysicalEvidenceSnapshot { - version: self.evidence_version.into(), - scope: self.target.scope.resolve()?, - cache_profile: self.target.scope.cache_profile.clone(), - storage_io: self.storage_io.cloned(), - handoffs: self.handoffs.cloned(), - }) - } - - fn query_node_evidence( - &self, - snapshot: &PhysicalEvidenceSnapshot, - request: PhysicalNodeRequest<'_>, - ) -> Result { - if snapshot.scope != self.target.scope.resolve()? { - return Err(AnalyticalCostError::ComparisonScopeMismatch( - "planner evidence snapshot", - )); - } - let mut matches = self - .candidate - .query_nodes() - .iter() - .enumerate() - .filter(|(_, entry)| { - entry.logical_node == *request.logical_node - && entry.operator == request.operator - && entry.occurrence == request.occurrence - && entry.synthetic == request.synthetic - }); - let (index, evidence) = matches.next().ok_or_else(|| { - AnalyticalCostError::MissingOperatorStatistics(format!( - "logical occurrence {}", - request.occurrence - )) - })?; - if matches.next().is_some() { - return Err(AnalyticalCostError::InvalidPhysicalDAG( - "duplicate query-node evidence key", - )); - } - self.used_query_nodes.borrow_mut().insert(index); - Ok(evidence.evidence.clone()) - } - - fn summary_physical_dag( - &self, - snapshot: &PhysicalEvidenceSnapshot, - _summary: &Rc, - _target: &asap_logical_optimizer::pass1::replacement::TargetSubDAG<'_>, - ) -> Result { - if snapshot.scope != self.target.scope.resolve()? { - return Err(AnalyticalCostError::ComparisonScopeMismatch( - "planner evidence snapshot", - )); - } - self.candidate - .summary_dag() - .cloned() - .ok_or(AnalyticalCostError::InvalidPhysicalDAG( - "summary candidate is missing its physical DAG", - )) - } -} - -struct ExportPlannerCostModel<'a> { - document: &'a PlannerCostDocument, -} - -impl ExportPlannerCostModel<'_> { - fn bound<'a>( - &'a self, - candidate: &ReplacementSubDAG, - target: &asap_logical_optimizer::pass1::replacement::TargetSubDAG<'_>, - ) -> Option<(ExportPhysicalProvider<'a>, &'a ResourceCalibration)> { - let mut targets = self - .document - .targets - .iter() - .filter(|entry| entry.target == *target.root); - let target_evidence = targets.next()?; - if targets.next().is_some() { - return None; - } - let mut candidates = target_evidence - .candidates - .iter() - .filter(|entry| entry.matches(candidate)); - let candidate_evidence = candidates.next()?; - if candidates.next().is_some() { - return None; - } - Some(( - ExportPhysicalProvider { - storage_io: self.document.storage_io.as_ref(), - handoffs: self.document.handoffs.as_ref(), - evidence_version: &self.document.evidence_version, - target: target_evidence, - candidate: candidate_evidence, - used_query_nodes: RefCell::new(std::collections::HashSet::new()), - }, - &self.document.calibration, - )) - } - - fn annotations( - &self, - candidate: &ReplacementSubDAG, - target: &Rc, - ) -> (CostAnnotation, CostAnnotation, CostAnnotation) { - let target = asap_logical_optimizer::pass1::replacement::TargetSubDAG::new(target); - let Some((provider, calibration)) = self.bound(candidate, &target) else { - return winner_cost_annotations(); - }; - let Ok(model) = PhysicalPlanCostModel::new(&provider, calibration.clone()) else { - return winner_cost_annotations(); - }; - let Ok(estimate) = model.estimate_candidate(candidate, &target) else { - return winner_cost_annotations(); - }; - if !provider.all_query_evidence_used() { - return winner_cost_annotations(); - } - let mut version = format!("{}+{}", ANALYTICAL_COST_MODEL_VERSION, calibration.version); - if let Some((storage, _)) = &estimate.storage_io { - version.push_str(&format!( - "+{}+{}", - storage.model_version, storage.calibration_version - )); - } - if let Some((handoff, _)) = &estimate.handoffs { - version.push_str(&format!( - "+{}+{}", - handoff.model_version, handoff.calibration_version - )); - } - let scope = &provider.target.scope; - let Ok((result_hits, buffer_hits)) = cache_hit_ratios( - &scope.cache_profile, - scope.evaluation_count, - scope.data_arrival, - ) else { - return winner_cost_annotations(); - }; - let input = |name: &str, value: f64, unit: &str| CostInput { - name: name.into(), - value, - unit: Some(unit.into()), - }; - let mut cache_inputs = vec![ - input( - "evaluation_count", - scope.evaluation_count as f64, - "evaluations", - ), - input("result_cache_hit_ratio", result_hits, "ratio"), - input("buffer_cache_hit_ratio", buffer_hits, "ratio"), - ]; - if let CacheProfile::Evidence(evidence) = &scope.cache_profile { - cache_inputs.extend([ - input( - "distinct_evaluations", - evidence.distinct_evaluations as f64, - "evaluations", - ), - input( - "repeated_identical_evaluations", - evidence.repeated_identical_evaluations as f64, - "evaluations", - ), - input( - "result_cache_working_set", - evidence.result_cache.working_set_bytes as f64, - "bytes", - ), - input( - "result_cache_capacity", - evidence.result_cache.capacity_bytes as f64, - "bytes", - ), - input( - "buffer_cache_working_set", - evidence.buffer_cache.working_set_bytes as f64, - "bytes", - ), - input( - "buffer_cache_capacity", - evidence.buffer_cache.capacity_bytes as f64, - "bytes", - ), - ]); - if let Some(ratio) = evidence.result_invalidation_ratio { - cache_inputs.push(input("result_cache_invalidation_ratio", ratio, "ratio")); - } - } - let inputs = |resources: asap_plan_selection::cost::analytical_cost::ResourceEstimate| { - let mut inputs = vec![ - CostInput { - name: "estimated_cpu_ops".into(), - value: resources.cpu_ops(), - unit: Some("operations".into()), - }, - CostInput { - name: "estimated_peak_memory".into(), - value: resources.peak_memory_bytes() as f64, - unit: Some("bytes".into()), - }, - CostInput { - name: "estimated_scan".into(), - value: resources.scan_bytes() as f64, - unit: Some("bytes".into()), - }, - ]; - inputs.extend(cache_inputs.iter().cloned()); - inputs - }; - let storage_inputs = |storage: &asap_plan_selection::cost::storage_io::StorageEstimate| { - let mut terms: Vec<_> = storage - .total - .terms() - .into_iter() - .map(|(name, value)| CostInput { - name: name.into(), - value: value as f64, - unit: Some("operations".into()), - }) - .collect(); - let mut ids: Vec<_> = storage.per_node.keys().collect(); - ids.sort(); - for id in ids { - terms.extend( - storage.per_node[id] - .terms() - .into_iter() - .map(|(name, value)| CostInput { - name: format!("physical_node:{id}:{name}"), - value: value as f64, - unit: Some("operations".into()), - }), - ); - } - terms - }; - let mut raw_inputs = inputs(estimate.resources.raw); - let mut candidate_inputs = inputs(estimate.resources.candidate); - if let Some((raw, candidate)) = &estimate.storage_io { - raw_inputs.extend(storage_inputs(raw)); - candidate_inputs.extend(storage_inputs(candidate)); - } - let handoff_inputs = - |estimate: &asap_plan_selection::cost::physical_handoff_cost::PhysicalHandoffEstimate| { - let mut terms: Vec<_> = estimate - .total - .terms() - .into_iter() - .map(|(name, value)| CostInput { - name: name.into(), - value: value as f64, - unit: Some("bytes".into()), - }) - .collect(); - for (prefix, entries) in [ - ("physical_node", &estimate.per_node), - ("boundary", &estimate.per_handoff), - ] { - let mut ids: Vec<_> = entries.keys().collect(); - ids.sort(); - for id in ids { - terms.extend(entries[id].terms().into_iter().map(|(name, value)| { - CostInput { - name: format!("{prefix}:{id}:{name}"), - value: value as f64, - unit: Some("bytes".into()), - } - })); - } - } - terms - }; - if let Some((raw, candidate)) = &estimate.handoffs { - raw_inputs.extend(handoff_inputs(raw)); - candidate_inputs.extend(handoff_inputs(candidate)); - } - let baseline = CostAnnotation::modeled( - estimate.raw_cost.0, - CostUnit::CostUnits, - &version, - raw_inputs, - ) - .with_evidence_version(&self.document.evidence_version) - .with_cache_profile(snapshot_cache_version(&provider)); - let selected = CostAnnotation::modeled( - estimate.candidate_cost.0, - CostUnit::CostUnits, - &version, - candidate_inputs, - ) - .with_baseline(BaselineRef::PreAsapRecomputation, estimate.raw_cost.0) - .with_evidence_version(&self.document.evidence_version) - .with_cache_profile(snapshot_cache_version(&provider)); - let benefit = CostAnnotation { - value: selected.delta, - unit: CostUnit::CostUnits, - source: CostSource::Modeled, - baseline: Some(BaselineRef::PreAsapRecomputation), - delta: None, - benefit_ratio: selected.benefit_ratio, - model_version: Some(version), - evidence_version: Some(self.document.evidence_version.clone()), - cache_profile: Some(snapshot_cache_version(&provider).into()), - benchmark_id: None, - inputs: Vec::new(), - }; - (baseline, selected, benefit) - } -} - -fn snapshot_cache_version<'a>(provider: &'a ExportPhysicalProvider<'_>) -> &'a str { - provider.target.scope.cache_profile.version() -} - -impl CostModel for ExportPlannerCostModel<'_> { - fn candidate_cost_covers_complete_plan(&self) -> bool { - true - } - - fn candidate_cost( - &self, - candidate: &ReplacementSubDAG, - target: &asap_logical_optimizer::pass1::replacement::TargetSubDAG<'_>, - ) -> Option { - let (provider, calibration) = self.bound(candidate, target)?; - let cost = PhysicalPlanCostModel::new(&provider, calibration.clone()) - .ok()? - .candidate_cost(candidate, target)?; - provider.all_query_evidence_used().then_some(cost) - } - - fn rank_candidates( - &self, - _intent: &asap_types::ir::operator::AggIntent, - candidates: &[asap_types::ir::schema::SketchAlgorithm], - ) -> Vec { - // Candidate generation must not reintroduce the legacy structural - // cost model before complete physical alternatives are compared. - candidates.to_vec() - } - - fn estimate_cost( - &self, - candidate: &ReplacementSubDAG, - target: &asap_logical_optimizer::pass1::replacement::TargetSubDAG<'_>, - ) -> f64 { - self.candidate_cost(candidate, target) - .map_or(f64::NAN, |cost| cost.0) - } -} - -/// Baseline/selected/benefit [`CostAnnotation`]s for one [`Winner`] — issue -/// #286's "replacement-region baseline cost, selected cost, and benefit" -/// granularity item, reused verbatim for [`TargetReplacement`] and for the -/// [`DAGDecision`] carried by every node the winning candidate produced or -/// carried. -/// -/// `dag_export` has no deployment-owned physical evidence provider. It must -/// therefore expose costs as unavailable instead of guessing operator -/// statistics or falling back to structural node counts. Callers that have -/// complete evidence use `PhysicalPlanCostModel` before export and may -/// attach its dimensional comparison to these fields. -#[allow(dead_code)] -fn winner_cost_annotations() -> (CostAnnotation, CostAnnotation, CostAnnotation) { - ( - CostAnnotation::unavailable(CostUnit::CostUnits), - CostAnnotation::unavailable(CostUnit::CostUnits), - CostAnnotation::unavailable(CostUnit::CostUnits), - ) -} - -use asap_devtools::{lower_promql_with_data_ingestion_interval, lower_sql, SqlCatalog}; - -enum Lang { - Sql, - PromQl, -} - -fn default_catalog() -> SqlCatalog { - SqlCatalog::new() - .with_table( - "metrics", - Schema::with_time_index( - vec![ - Field::plain("ts", DataType::Timestamp, false), - Field::plain("service", DataType::Utf8, false), - Field::plain("region", DataType::Utf8, false), - Field::plain("latency", DataType::Float64, false), - Field::plain("bytes", DataType::Int64, false), - ], - 0, - vec![], - ), - ) - .with_table( - "hosts", - Schema::new(vec![ - Field::plain("service", DataType::Utf8, false), - Field::plain("region", DataType::Utf8, false), - ]), - ) -} - -fn catalog(custom: &[String]) -> SqlCatalog { - let mut catalog = default_catalog(); - for raw in custom { - let value: serde_json::Value = serde_json::from_str(raw) - .unwrap_or_else(|error| panic!("--table-schema must be valid JSON: {error}")); - let name = value["name"] - .as_str() - .expect("--table-schema.name must be a string"); - let columns = value["columns"] - .as_array() - .expect("--table-schema.columns must be an array"); - let columns: Vec = columns - .iter() - .map(|column| { - let column_name = column["name"] - .as_str() - .expect("column.name must be a string"); - let data_type = match column["type"] - .as_str() - .expect("column.type must be a string") - .to_ascii_lowercase() - .as_str() - { - "timestamp" => DataType::Timestamp, - "utf8" | "string" => DataType::Utf8, - "float64" | "double" => DataType::Float64, - "int64" | "bigint" => DataType::Int64, - other => panic!("unsupported column type {other:?}"), - }; - Field::plain( - column_name, - data_type, - column["nullable"].as_bool().unwrap_or(true), - ) - }) - .collect(); - let schema = match value.get("time_index").and_then(|index| index.as_u64()) { - Some(index) => Schema::with_time_index(columns, index as usize, vec![]), - None => Schema::new(columns), - }; - catalog = catalog.with_table(name, schema); - } - catalog -} - -/// Parses `--sql "" --name "