Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
117 changes: 93 additions & 24 deletions crates/devtools/src/bin/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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 '<json>'` (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,
};
Expand All @@ -80,8 +90,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 <query>... \
[--epsilon <f64> --delta <f64>] [--interval-ms <u64>]) [--max-candidates <n>] --out <file>";
"usage: stage_pipeline (--example planner-layering-{1,2,3a,3b,4a,4b} | (--promql <query>... \
| --table <json>... --sql <query>...) [--epsilon <f64> --delta <f64>] [--interval-ms <u64>]) \
[--max-candidates <n>] --out <file>";

fn main() {
if let Err(message) = run(std::env::args().skip(1).collect()) {
Expand All @@ -92,7 +103,7 @@ fn main() {

fn run(args: Vec<String>) -> 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;
Expand All @@ -102,7 +113,9 @@ fn run(args: Vec<String>) -> Result<(), String> {
let number = |v: String| v.parse::<f64>().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}"))?,
Expand All @@ -111,9 +124,22 @@ fn run(args: Vec<String>) -> 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(),
Expand All @@ -126,17 +152,18 @@ fn run(args: Vec<String>) -> 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}"))
}
Expand All @@ -147,17 +174,17 @@ const PRICE_LIMIT: usize = 4096;

fn stage_pipeline(
workload: &PlanningWorkload,
catalog: &SqlCatalog,
capabilities: &DeploymentCapabilities,
max_candidates: usize,
) -> Result<Value, String> {
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| {
Expand Down Expand Up @@ -502,14 +529,15 @@ fn declared<T>(value: T) -> Evidence<T> {
}
}

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()
Expand Down Expand Up @@ -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<SqlCatalog, String> {
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::<Result<Vec<_>, _>>()?;
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
Expand Down
58 changes: 56 additions & 2 deletions crates/devtools/tests/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -319,3 +318,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}");
}
18 changes: 14 additions & 4 deletions tools/dag-viewer/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -44,10 +44,20 @@ the page, or with `?doc=<path>` 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 <json> … --sql …` and the page shows the result.

## Document format

Expand Down
42 changes: 36 additions & 6 deletions tools/dag-viewer/editor.js
Original file line number Diff line number Diff line change
@@ -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);

Expand All @@ -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;
Expand All @@ -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}`;
Expand Down
Loading