diff --git a/crates/devtools/src/bin/stage_pipeline.rs b/crates/devtools/src/bin/stage_pipeline.rs index be6c47a0..f66845aa 100644 --- a/crates/devtools/src/bin/stage_pipeline.rs +++ b/crates/devtools/src/bin/stage_pipeline.rs @@ -30,6 +30,10 @@ // The deployment inputs are the built-in cost and accuracy models and the // reference executor's capabilities (`asap_executor::capabilities`). // +// - deployment: the deployment inputs Stage 3 used (the executor's +// capabilities, summarized; the cost model and its Stage3Calibration; +// the accuracy model's name). +// // Everything is the library's `plan_selection::plan_stages`, the function the // facade runs; this tool only serializes it. Stage 3 here is over every // displayed candidate; the facade's dynamic program selects the same winner @@ -46,7 +50,9 @@ use asap_logical_optimizer::pass1::logical_candidates::{ }; use asap_logical_optimizer::Realization; use asap_plan_selection::PlanningModels; -use asap_plan_selection::{plan_stages, Selection, Sharing, MAX_ENUMERATED_CANDIDATES}; +use asap_plan_selection::{ + plan_stages, Selection, Sharing, COST_MODEL, COST_PER_SECOND, MAX_ENUMERATED_CANDIDATES, +}; use asap_types::ir::export::{ compile_logical_asap_workload, LogicalASAPDAG, LogicalASAPDAGDocument, }; @@ -138,11 +144,12 @@ fn stage_pipeline(workload: &PlanningWorkload, max_candidates: usize) -> Result< .collect(); let data = workload.data_workload.clone().unwrap_or_default(); let capabilities = asap_executor::capabilities(); + let models = PlanningModels::builtin().with_capabilities(&capabilities); let run = plan_stages( roots.into_iter().enumerate().collect(), &demand, &data, - PlanningModels::builtin().with_capabilities(&capabilities), + models, max_candidates.max(1), ) .map_err(|e| format!("planning: {e}"))?; @@ -190,6 +197,7 @@ fn stage_pipeline(workload: &PlanningWorkload, max_candidates: usize) -> Result< Ok(json!({ "format": "asap-stage-pipeline/v1", "workload": { "queries": workload_queries(workload) }, + "deployment": deployment_json(&models), "stage0_logical": { "dag": stage0 }, "stage1_logical_asap": { "combinations": combinations, @@ -201,6 +209,47 @@ fn stage_pipeline(workload: &PlanningWorkload, max_candidates: usize) -> Result< })) } +/// The deployment inputs Stage 3 used: the executor's capabilities, one +/// line per summary, and the cost and accuracy models. +fn deployment_json(models: &PlanningModels<'_>) -> Value { + let caps = models.capabilities; + let summaries = caps.summaries.as_ref().map(|summaries| { + summaries + .iter() + .map(|s| { + let readouts: Vec<_> = s.readouts.iter().map(|r| format!("{r:?}")).collect(); + json!({ "summary": s.name(), "readouts": readouts }) + }) + .collect::>() + }); + let c = &models.calibration; + json!({ + "capabilities": { + "source": "asap_executor::capabilities", + "summaries": summaries, + "ingestion_time": caps.ingestion_time, + "query_time_retention": caps.query_time_retention, + "memory_budget_bytes": caps.memory_budget_bytes, + "raw_data_retained": caps.raw_data_retained, + "raw_bytes_per_sample": caps.raw_bytes_per_sample, + }, + "cost_model": { + "name": COST_MODEL, + "unit": COST_PER_SECOND, + "calibration": { + "version": c.version, + "cost_per_cpu_op": c.cost_per_cpu_op, + "cost_per_scan_byte": c.cost_per_scan_byte, + "cost_per_retained_byte_second": c.cost_per_retained_byte_second, + "horizon_s": c.horizon_s, + "latency_ms_per_cost_unit": c.latency_ms_per_cost_unit, + }, + }, + // `PlanningModels::builtin()`'s accuracy model and evidence. + "accuracy_model": { "name": "DefaultAccuracyModel", "evidence": "none" }, + }) +} + fn stage3_json(selection: &Selection) -> Value { let costs: serde_json::Map<_, _> = selection .costs diff --git a/crates/devtools/tests/stage_pipeline.rs b/crates/devtools/tests/stage_pipeline.rs index 8c5a1815..824df584 100644 --- a/crates/devtools/tests/stage_pipeline.rs +++ b/crates/devtools/tests/stage_pipeline.rs @@ -14,9 +14,12 @@ const COMMITTED: &str = "../../tools/dag-viewer/examples/planner-layering-exampl const EXAMPLE1: [&str; 4] = ["--example", "planner-layering-1", "--max-candidates", "128"]; fn generate(args: &[&str]) -> Value { + // Tests run in parallel and may generate the same document. + static NEXT: std::sync::atomic::AtomicUsize = std::sync::atomic::AtomicUsize::new(0); let out = std::env::temp_dir().join(format!( - "stage_pipeline_{}_{}.json", + "stage_pipeline_{}_{}_{}.json", std::process::id(), + NEXT.fetch_add(1, std::sync::atomic::Ordering::Relaxed), args.join("_") )); let status = Command::new(env!("CARGO_BIN_EXE_stage_pipeline")) @@ -191,3 +194,33 @@ fn example4a_repeats_monthly_with_nothing_maintainable() { }; assert!(cost(&monthly) < cost(&once)); } + +/// The document records the deployment inputs Stage 3 used: the executor's +/// capabilities (Hydra summaries named by their kind), the cost model with +/// its calibration, including the memory weight and raw-retention settings, +/// and the accuracy model's name. +#[test] +fn document_records_the_deployment_inputs() { + let document = generate(&EXAMPLE1); + let deployment = &document["deployment"]; + let capabilities = &deployment["capabilities"]; + assert_eq!(capabilities["source"], "asap_executor::capabilities"); + assert_eq!(capabilities["ingestion_time"], true); + assert_eq!(capabilities["raw_data_retained"], true); + assert_eq!(capabilities["raw_bytes_per_sample"], 16); + assert!(capabilities["memory_budget_bytes"].is_null()); + let summaries = capabilities["summaries"].as_array().unwrap(); + let hydra = summaries + .iter() + .find(|s| s["summary"] == "HydraCms") + .expect("HydraCms is listed"); + assert_eq!( + hydra["readouts"], + serde_json::json!(["TotalCount", "ItemCount"]) + ); + let calibration = &deployment["cost_model"]["calibration"]; + assert_eq!(deployment["cost_model"]["name"], "analytical-cost-v2"); + assert_eq!(calibration["cost_per_retained_byte_second"], 1.25e-7); + assert_eq!(calibration["version"], "illustrative-v2"); + assert_eq!(deployment["accuracy_model"]["name"], "DefaultAccuracyModel"); +} diff --git a/crates/plan-selection/src/lib.rs b/crates/plan-selection/src/lib.rs index da99b22b..9f0448d7 100644 --- a/crates/plan-selection/src/lib.rs +++ b/crates/plan-selection/src/lib.rs @@ -100,6 +100,9 @@ use asap_physical_optimizer::materialization::MAX_PHYSICAL_PER_LOGICAL; /// [`Stage3Calibration::ILLUSTRATIVE`]. See /// `docs/design_docs/proposals/stage3-cost-model.md`. pub const COST_PER_SECOND: &str = "cost_per_second"; +/// The analytical cost model Stage 3 prices with, as reported in each +/// candidate's cost `source`. +pub const COST_MODEL: &str = "analytical-cost-v2"; /// Groups assumed for every `by (...)` reduction, absent group-count evidence. const DEFAULT_GROUP_COUNT: u64 = 100; @@ -1877,7 +1880,7 @@ fn price_nodes( total: per_node.values().map(|n| n.cost).sum(), unit: COST_PER_SECOND, source: format!( - "analytical-cost-v2 (illustrative statistics, calibration {})", + "{COST_MODEL} (illustrative statistics, calibration {})", calibration.version ), per_node, diff --git a/crates/types/src/deployment.rs b/crates/types/src/deployment.rs index d157e58a..64e4ce7c 100644 --- a/crates/types/src/deployment.rs +++ b/crates/types/src/deployment.rs @@ -101,6 +101,13 @@ pub struct SummarySupport { pub readouts: BTreeSet, } +impl SummarySupport { + /// E.g. "Kll", "HydraCms" or "exact Sum". + pub fn name(&self) -> String { + name(&self.family, &self.layout) + } +} + /// A summary family, without its sizing parameters. #[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] pub enum SummaryFamily { diff --git a/tools/dag-viewer/examples/planner-layering-example1.json b/tools/dag-viewer/examples/planner-layering-example1.json index f8995b86..d3e2b788 100644 --- a/tools/dag-viewer/examples/planner-layering-example1.json +++ b/tools/dag-viewer/examples/planner-layering-example1.json @@ -1,4 +1,104 @@ { + "deployment": { + "accuracy_model": { + "evidence": "none", + "name": "DefaultAccuracyModel" + }, + "capabilities": { + "ingestion_time": true, + "memory_budget_bytes": null, + "query_time_retention": false, + "raw_bytes_per_sample": 16, + "raw_data_retained": true, + "source": "asap_executor::capabilities", + "summaries": [ + { + "readouts": [], + "summary": "exact Sum" + }, + { + "readouts": [], + "summary": "exact Count" + }, + { + "readouts": [], + "summary": "exact Min" + }, + { + "readouts": [], + "summary": "exact Max" + }, + { + "readouts": [], + "summary": "exact Increase" + }, + { + "readouts": [], + "summary": "exact Rate" + }, + { + "readouts": [ + "Quantile" + ], + "summary": "Kll" + }, + { + "readouts": [ + "Quantile", + "TotalCount" + ], + "summary": "DDSketch" + }, + { + "readouts": [ + "TotalCount", + "Cardinality" + ], + "summary": "Hll" + }, + { + "readouts": [ + "TotalCount", + "ItemCount" + ], + "summary": "HydraCms" + }, + { + "readouts": [ + "TopK" + ], + "summary": "CmsWithHeap" + }, + { + "readouts": [ + "TopK" + ], + "summary": "CountSketchWithHeap" + }, + { + "readouts": [ + "TotalCount", + "Cardinality", + "FrequencyL2", + "FrequencyEntropy" + ], + "summary": "UnivMon" + } + ] + }, + "cost_model": { + "calibration": { + "cost_per_cpu_op": 1e-6, + "cost_per_retained_byte_second": 1.25e-7, + "cost_per_scan_byte": 1e-7, + "horizon_s": 3600.0, + "latency_ms_per_cost_unit": 1.0, + "version": "illustrative-v2" + }, + "name": "analytical-cost-v2", + "unit": "cost_per_second" + } + }, "format": "asap-stage-pipeline/v1", "stage0_logical": { "dag": {