diff --git a/crates/devtools/src/bin/stage_pipeline.rs b/crates/devtools/src/bin/stage_pipeline.rs index dd4ccdd2..a5a597c0 100644 --- a/crates/devtools/src/bin/stage_pipeline.rs +++ b/crates/devtools/src/bin/stage_pipeline.rs @@ -7,6 +7,10 @@ // every 10 min on a deployment that does not keep raw data) // cargo run -p asap-devtools --bin stage_pipeline -- \ // --promql "topk by (job) (10, rate(x[1m]))" --epsilon 0.01 --delta 0.001 --out run.json +// cargo run -p asap-devtools --bin stage_pipeline -- \ +// --table '{"name": "flows", "columns": [{"name": "ts", "type": "timestamp"}, +// {"name": "src_ip", "type": "utf8", "nullable": false}], "time_index": 0}' \ +// --sql "SELECT COUNT(DISTINCT src_ip) FROM flows" --epsilon 0.02 --out run.json // // Writes an `asap-stage-pipeline/v1` document (tools/dag-viewer) with the // four planner stages (#509 MVP): @@ -48,13 +52,19 @@ // displayed candidate; the facade's dynamic program selects the same winner // when its assumptions hold. // -// `--promql` may repeat. `--epsilon`/`--delta` apply to every `--promql` -// query; without them the queries are exact. `--interval-ms` is the source -// cadence PromQL needs (default 15000). +// `--promql` and `--sql` may repeat; a run is one language. `--epsilon`/ +// `--delta` apply to every query; without them the queries are exact. +// `--interval-ms` is the source cadence PromQL needs (default 15000). +// `--table ''` (repeatable) declares a table SQL queries read: +// `{"name": ..., "columns": [{"name": ..., "type": "timestamp|utf8|string| +// float64|double|int64|bigint", "nullable": true}], "time_index": 0}` +// (`nullable` defaults to true; without `time_index` the table has no time +// column). There are no default tables. use std::collections::HashSet; use std::rc::Rc; +use asap_frontend_sql::SqlCatalog; use asap_logical_optimizer::pass1::logical_candidates::{ choice_index, combination_count, LocalLogicalCandidates, }; @@ -78,8 +88,9 @@ use asap_types::workload::{ use serde_json::{json, Value}; const USAGE: &str = - "usage: stage_pipeline (--example planner-layering-{1,2,3a,3b,4a,4b} | --promql ... \ -[--epsilon --delta ] [--interval-ms ]) [--max-candidates ] --out "; + "usage: stage_pipeline (--example planner-layering-{1,2,3a,3b,4a,4b} | (--promql ... \ +| --table ... --sql ...) [--epsilon --delta ] [--interval-ms ]) \ +[--max-candidates ] --out "; fn main() { if let Err(message) = run(std::env::args().skip(1).collect()) { @@ -90,7 +101,7 @@ fn main() { fn run(args: Vec) -> Result<(), String> { let mut example = None; - let mut queries = Vec::new(); + let (mut promql, mut sql, mut tables) = (Vec::new(), Vec::new(), Vec::new()); let (mut epsilon, mut delta, mut interval_ms) = (None, None, 15_000u64); let mut max_candidates = MAX_ENUMERATED_CANDIDATES; let mut out = None; @@ -100,7 +111,9 @@ fn run(args: Vec) -> Result<(), String> { let number = |v: String| v.parse::().map_err(|e| format!("{v}: {e}")); match flag.as_str() { "--example" => example = Some(value()?), - "--promql" => queries.push(value()?), + "--promql" => promql.push(value()?), + "--sql" => sql.push(value()?), + "--table" => tables.push(value()?), "--epsilon" => epsilon = Some(number(value()?)?), "--delta" => delta = Some(number(value()?)?), "--interval-ms" => interval_ms = value()?.parse().map_err(|e| format!("{e}"))?, @@ -109,9 +122,22 @@ fn run(args: Vec) -> Result<(), String> { other => return Err(format!("unknown argument {other}")), } } + let (language, queries) = match (promql.is_empty(), sql.is_empty()) { + (false, false) => { + return Err("a run is one language: give --promql or --sql, not both".into()) + } + (true, false) => (QueryLanguage::SQL(SqlDialect::DataFusionSQL), sql), + _ => (QueryLanguage::PromQL, promql), + }; + if !tables.is_empty() && !matches!(language, QueryLanguage::SQL(_)) { + return Err("--table needs --sql".into()); + } let workload = match (example.as_deref(), queries.is_empty()) { (Some("planner-layering-1"), true) => planner_layering_example1(), - (Some("planner-layering-2"), true) => planner_layering_example2(), + (Some("planner-layering-2"), true) => { + tables = vec![FLOWS.into()]; + planner_layering_example2() + } (Some("planner-layering-3a"), true) => planner_layering_example3a(), (Some("planner-layering-3b"), true) => planner_layering_example3b(), (Some("planner-layering-4a"), true) => planner_layering_example4a(), @@ -124,17 +150,18 @@ fn run(args: Vec) -> Result<(), String> { (Some(epsilon), Some(delta)) => AccuracyTarget::EpsilonDelta { epsilon, delta }, (None, Some(_)) => return Err("--delta needs --epsilon".into()), }; - promql_batch(&queries, accuracy, interval_ms) + query_batch(language, &queries, accuracy, interval_ms) } - _ => return Err("give exactly one of --example or --promql".into()), + _ => return Err("give exactly one of --example, --promql or --sql".into()), }; + let catalog = sql_catalog(&tables)?; let out = out.ok_or("--out is required")?; let capabilities = DeploymentCapabilities { // Example 4b's crossover needs a deployment that does not keep raw data. raw_data_retained: example.as_deref() != Some("planner-layering-4b"), ..asap_executor::capabilities() }; - let document = stage_pipeline(&workload, &capabilities, max_candidates)?; + let document = stage_pipeline(&workload, &catalog, &capabilities, max_candidates)?; let text = serde_json::to_string_pretty(&document).map_err(|e| e.to_string())? + "\n"; std::fs::write(&out, text).map_err(|e| format!("{out}: {e}")) } @@ -145,17 +172,17 @@ const PRICE_LIMIT: usize = 4096; fn stage_pipeline( workload: &PlanningWorkload, + catalog: &SqlCatalog, capabilities: &DeploymentCapabilities, max_candidates: usize, ) -> Result { let roots = match workload.query_workload.language { - // The only SQL workload here is Example 2's, over `flows`. QueryLanguage::SQL(_) => tokio::runtime::Builder::new_current_thread() .build() .map_err(|e| e.to_string())? .block_on(asap_frontend_sql::lower_sql_batch( &workload.query_workload, - &flows_catalog(), + catalog, )) .into_iter() .map(|root| { @@ -502,14 +529,15 @@ fn declared(value: T) -> Evidence { } } -fn promql_batch( +fn query_batch( + language: QueryLanguage, queries: &[String], accuracy: AccuracyTarget, interval_ms: u64, ) -> PlanningWorkload { PlanningWorkload { query_workload: QueryWorkload { - language: QueryLanguage::PromQL, + language, query_batch: Some( queries .iter() @@ -585,15 +613,56 @@ fn planner_layering_example1() -> PlanningWorkload { } } -/// #509 Example 2's `flows` table. -fn flows_catalog() -> asap_frontend_sql::SqlCatalog { - asap_frontend_sql::SqlCatalog::new().with_table( - "flows", - Schema::new(vec![ - Field::plain("ts", DataType::Timestamp, false), - Field::plain("src_ip", DataType::Utf8, false), - ]), - ) +/// #509 Example 2's `flows` table, as a `--table`. +const FLOWS: &str = r#"{"name": "flows", "columns": [ + {"name": "ts", "type": "timestamp", "nullable": false}, + {"name": "src_ip", "type": "utf8", "nullable": false}]}"#; + +/// The SQL catalog of the `--table` declarations. +fn sql_catalog(tables: &[String]) -> Result { + let mut catalog = SqlCatalog::new(); + let mut names = HashSet::new(); + for raw in tables { + let bad = |what: &str| format!("--table {raw}: {what}"); + let value: Value = serde_json::from_str(raw).map_err(|e| bad(&e.to_string()))?; + let name = value["name"] + .as_str() + .filter(|name| !name.is_empty()) + .ok_or_else(|| bad("name must be a non-empty string"))?; + if !names.insert(name.to_string()) { + return Err(bad("the table is declared twice")); + } + let columns = value["columns"] + .as_array() + .filter(|columns| !columns.is_empty()) + .ok_or_else(|| bad("columns must be a non-empty array"))? + .iter() + .map(|column| { + let column_name = column["name"] + .as_str() + .filter(|name| !name.is_empty()) + .ok_or_else(|| bad("column.name must be a non-empty string"))?; + let data_type = match column["type"].as_str().map(str::to_ascii_lowercase) { + Some(t) if t == "timestamp" => DataType::Timestamp, + Some(t) if t == "utf8" || t == "string" => DataType::Utf8, + Some(t) if t == "float64" || t == "double" => DataType::Float64, + Some(t) if t == "int64" || t == "bigint" => DataType::Int64, + _ => return Err(bad(&format!("unsupported column type {}", column["type"]))), + }; + let nullable = column["nullable"].as_bool().unwrap_or(true); + Ok(Field::plain(column_name, data_type, nullable)) + }) + .collect::, _>>()?; + let schema = match value.get("time_index") { + None | Some(Value::Null) => Schema::new(columns), + Some(index) => match index.as_u64().map(|i| i as usize) { + Some(i) if i < columns.len() => Schema::with_time_index(columns, i, vec![]), + _ => return Err(bad("time_index must be a column index")), + }, + }; + catalog = catalog.with_table(name, schema); + } + Ok(catalog) } /// #509 Example 2: distinct count, entropy and L2 of `src_ip` over the last diff --git a/crates/devtools/tests/stage_pipeline.rs b/crates/devtools/tests/stage_pipeline.rs index 6041ede3..cfabec6f 100644 --- a/crates/devtools/tests/stage_pipeline.rs +++ b/crates/devtools/tests/stage_pipeline.rs @@ -17,10 +17,9 @@ 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")) .args(args) @@ -332,3 +331,58 @@ fn document_records_the_deployment_inputs() { assert_eq!(calibration["version"], "illustrative-v2"); assert_eq!(deployment["accuracy_model"]["name"], "DefaultAccuracyModel"); } + +const FLOWS: &str = r#"{"name": "flows", "columns": [{"name": "ts", "type": "timestamp", "nullable": false}, {"name": "src_ip", "type": "utf8", "nullable": false}], "time_index": 0}"#; + +/// A `--sql` query over a `--table` lowers through the SQL frontend, plans, +/// and is reported as SQL with its accuracy target. +#[test] +fn sql_over_declared_tables_plans() { + let document = generate(&[ + "--table", + FLOWS, + "--sql", + "SELECT COUNT(DISTINCT src_ip) FROM flows", + "--epsilon", + "0.02", + ]); + assert_valid_document(&document, 1); + let query = &document["workload"]["queries"][0]; + assert_eq!(query["language"], "sql"); + assert_eq!(query["requirements"]["accuracy"]["epsilon"], 0.02); +} + +/// The tool's error message for `args`, which must fail. +fn failure(args: &[&str]) -> String { + let output = Command::new(env!("CARGO_BIN_EXE_stage_pipeline")) + .args(args) + .args(["--out", "unused.json"]) + .output() + .unwrap(); + assert!(!output.status.success(), "{args:?} should fail"); + String::from_utf8(output.stderr).unwrap() +} + +/// A run is one language, `--table` needs `--sql`, and a malformed table +/// or a query over an undeclared table is an error, not a panic. +#[test] +fn bad_sql_arguments_are_rejected() { + let query = "SELECT COUNT(*) FROM flows"; + let error = failure(&["--promql", "up", "--table", FLOWS, "--sql", query]); + assert!(error.contains("one language"), "{error}"); + let error = failure(&["--table", FLOWS, "--promql", "up"]); + assert!(error.contains("--table needs --sql"), "{error}"); + for table in [ + "{not json", + r#"{"columns": [{"name": "ts", "type": "timestamp"}]}"#, + r#"{"name": "flows", "columns": [{"name": "ts", "type": "decimal"}]}"#, + r#"{"name": "flows", "columns": [{"name": "ts", "type": "timestamp"}], "time_index": 1}"#, + ] { + let error = failure(&["--table", table, "--sql", query]); + assert!(error.contains("--table"), "{table}: {error}"); + } + let error = failure(&["--table", FLOWS, "--table", FLOWS, "--sql", query]); + assert!(error.contains("declared twice"), "{error}"); + let error = failure(&["--sql", query]); + assert!(error.contains("lowering"), "{error}"); +} diff --git a/tools/dag-viewer/README.md b/tools/dag-viewer/README.md index 626d015f..6d911977 100644 --- a/tools/dag-viewer/README.md +++ b/tools/dag-viewer/README.md @@ -44,10 +44,20 @@ the page, or with `?doc=` for a file served next to `index.html`. ## Query editor -With the local server running, **Query editor** plans PromQL queries (one -per line) with an optional ε and δ for every query and the sample interval. -The server runs `stage_pipeline --promql … --out …` and the page shows the -result. SQL needs a catalog and is not in the editor yet. +With the local server running, **Query editor** plans PromQL or SQL queries +(one per line, one language per run) with an optional ε and δ for every query. +PromQL also takes the sample interval. SQL queries read the tables declared +in the tables box, one JSON object per line: + +```json +{"name": "flows", "columns": [{"name": "ts", "type": "timestamp", "nullable": false}, + {"name": "src_ip", "type": "utf8", "nullable": false}], "time_index": 0} +``` + +Column types are `timestamp`, `utf8`/`string`, `float64`/`double` and +`int64`/`bigint`; `nullable` defaults to true and `time_index` (the time +column's position) is optional. The server runs `stage_pipeline --promql …` +or `stage_pipeline --table … --sql …` and the page shows the result. ## Document format diff --git a/tools/dag-viewer/editor.js b/tools/dag-viewer/editor.js index 54c9cf7f..4855ac96 100644 --- a/tools/dag-viewer/editor.js +++ b/tools/dag-viewer/editor.js @@ -1,5 +1,6 @@ -// Query editor: plans PromQL queries through the local server's /api/plan -// (server.py runs `stage_pipeline`) and shows the resulting stage document. +// Query editor: plans PromQL or SQL queries through the local server's +// /api/plan (server.py runs `stage_pipeline`) and shows the resulting stage +// document. SQL queries read the tables declared one JSON object per line. (function () { const $ = (id) => document.getElementById(id); @@ -9,16 +10,45 @@ $('editorToggle').setAttribute('aria-pressed', String(open)); }); + // Each language keeps its own queries while the other is shown. + let language = 'promql'; + const drafts = { sql: 'SELECT COUNT(DISTINCT src_ip) FROM flows' }; + const LABELS = { promql: 'PromQL', sql: 'SQL' }; + const setLanguage = (next) => { + if (next === language) return; + drafts[language] = $('editorQueries').value; + $('editorQueries').value = drafts[next]; + language = next; + $('editorPromql').setAttribute('aria-pressed', String(next === 'promql')); + $('editorSql').setAttribute('aria-pressed', String(next === 'sql')); + $('editorTablesField').hidden = next !== 'sql'; + $('editorIntervalField').hidden = next === 'sql'; + }; + $('editorPromql').addEventListener('click', () => setLanguage('promql')); + $('editorSql').addEventListener('click', () => setLanguage('sql')); + + const lines = (id) => $(id).value.split('\n').map((line) => line.trim()).filter(Boolean); + let busy = false; $('editorForm').addEventListener('submit', async (event) => { event.preventDefault(); if (busy) return; - const queries = $('editorQueries').value.split('\n').map((q) => q.trim()).filter(Boolean); + const queries = lines('editorQueries'); if (!queries.length) { - $('editorStatus').textContent = 'Enter at least one PromQL query.'; + $('editorStatus').textContent = `Enter at least one ${LABELS[language]} query.`; return; } - const body = { queries, interval_ms: Number($('editorInterval').value) || 15000 }; + const body = { language, queries, interval_ms: Number($('editorInterval').value) || 15000 }; + if (language === 'sql') { + try { + body.tables = lines('editorTables').map((line, i) => { + try { return JSON.parse(line); } catch (err) { throw new Error(`table line ${i + 1}: ${err.message}`); } + }); + } catch (err) { + $('editorStatus').textContent = err.message; + return; + } + } if ($('editorEpsilon').value) body.epsilon = Number($('editorEpsilon').value); if ($('editorDelta').value) body.delta = Number($('editorDelta').value); busy = true; @@ -31,7 +61,7 @@ }); const result = await response.json().catch(() => ({ error: `HTTP ${response.status}` })); if (!response.ok) throw new Error(result.error || `HTTP ${response.status}`); - window.StageViewer.showDocument(result, `editor · ${queries.length} quer${queries.length === 1 ? 'y' : 'ies'}`); + window.StageViewer.showDocument(result, `editor · ${queries.length} ${LABELS[language]} quer${queries.length === 1 ? 'y' : 'ies'}`); $('editorStatus').textContent = 'Planned.'; } catch (err) { $('editorStatus').textContent = `Planning failed: ${err.message}`; diff --git a/tools/dag-viewer/index.html b/tools/dag-viewer/index.html index b773116a..92864db4 100644 --- a/tools/dag-viewer/index.html +++ b/tools/dag-viewer/index.html @@ -121,6 +121,8 @@ .editor .row { display: flex; flex-wrap: wrap; gap: 12px; align-items: center; } .editor input { font: inherit; width: 8em; color: var(--fg); background: var(--panel-2); border: 1px solid var(--line); border-radius: 6px; padding: 3px 6px; } .editor label { font-size: 13px; color: var(--muted); display: flex; gap: 6px; align-items: center; } +.editor label.stack { flex-direction: column; align-items: stretch; } +.editor [hidden] { display: none; } @@ -136,14 +138,20 @@

ASAPPlanner Stage Viewer