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
122 changes: 115 additions & 7 deletions crates/devtools/src/bin/stage_pipeline.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
// cargo run -p asap-devtools --bin stage_pipeline -- \
// --example planner-layering-1 --out planner-layering-example1.json
// --example planner-layering-1 --max-candidates 128 --out planner-layering-example1.json
// (also planner-layering-3a and planner-layering-3b: #509 Example 3,
// Patterns A and B)
// 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
//
Expand All @@ -12,7 +14,9 @@
// merges something, again with identical sub-DAGs shared ("· shared
// input"), and, when the summary-capability rule applies, again with
// one summary sized for its strictest consumer ("· shared summary"); in
// enumeration order and capped by `--max-candidates` (default 64);
// enumeration order and capped by `--max-candidates` (default 64). A
// repeating query's mergeable alternatives also come in tumbling panes
// (Pass 2's window-composition rule), e.g. "Q1 Kll · tumbling 1m panes";
// - stage2_physical_asap: one physical candidate per logical candidate
// (operator implementation only, everything at query time), no cost;
// - stage3_selection: per-candidate costs, the selected candidate, and
Expand Down Expand Up @@ -47,11 +51,12 @@ use asap_types::workload::{
AccuracyRequirement, BatchEntry, DataArrival, DataDistribution, DataWorkload, DurationMs,
Evidence, EvidenceSource, LatencyRequirement, MetricType, PlanningWorkload, Predictability,
Query, QueryLanguage, QueryRecurrence, QueryRequirements, QueryTimeScope, QueryWorkload, Rate,
RepeatedDemand, RepeatingEntry, RepetitionInterval, RootDemand, TimeSelection,
RepeatedDemand, RepeatingEntry, RepetitionInterval, RootDemand, TimeSelection, TimestampMs,
};
use serde_json::{json, Value};

const USAGE: &str = "usage: stage_pipeline (--example planner-layering-1 | --promql <query>... \
const USAGE: &str =
"usage: stage_pipeline (--example planner-layering-{1,3a,3b} | --promql <query>... \
[--epsilon <f64> --delta <f64>] [--interval-ms <u64>]) [--max-candidates <n>] --out <file>";

fn main() {
Expand Down Expand Up @@ -84,6 +89,8 @@ fn run(args: Vec<String>) -> Result<(), String> {
}
let workload = match (example.as_deref(), queries.is_empty()) {
(Some("planner-layering-1"), true) => planner_layering_example1(),
(Some("planner-layering-3a"), true) => planner_layering_example3a(),
(Some("planner-layering-3b"), true) => planner_layering_example3b(),
(Some(other), true) => return Err(format!("unknown example {other}")),
(None, false) => {
let accuracy = match (epsilon, delta) {
Expand Down Expand Up @@ -245,7 +252,8 @@ fn target_owners(inventory: &LocalLogicalCandidates<usize>) -> Vec<usize> {

/// E.g. "Q1 exact · Q2 CMS+heap"; exact accumulators are listed in
/// parentheses. A sketch that absorbs the aggregate beneath it is
/// "whole-expression".
/// "whole-expression"; one built in tumbling panes names them, e.g.
/// "Kll · tumbling 1m panes".
fn label(inventory: &LocalLogicalCandidates<usize>, owners: &[usize], choice: &[usize]) -> String {
(0..inventory.roots.len())
.map(|query| {
Expand All @@ -257,10 +265,13 @@ fn label(inventory: &LocalLogicalCandidates<usize>, owners: &[usize], choice: &[
if owner != query {
continue;
}
let window = target.windows[index]
.label()
.map_or(String::new(), |form| format!(" · {form}"));
match &target.alternatives[index] {
Realization::PassThrough => {}
Realization::ExactAggregate { kind, .. } => {
accumulators.push(format!("{kind:?} acc"))
accumulators.push(format!("{kind:?} acc{window}"))
}
Realization::Sketch(kind) => {
let name = match kind.algorithm() {
Expand All @@ -271,7 +282,7 @@ fn label(inventory: &LocalLogicalCandidates<usize>, owners: &[usize], choice: &[
// The sketch reads the inner aggregate's input and replaces it.
sketches.push(match target.absorbs[index] {
Some(_) => format!("whole-expression {name}"),
None => name,
None => format!("{name}{window}"),
});
}
other => sketches.push(format!("{other:?}")),
Expand Down Expand Up @@ -414,3 +425,100 @@ fn planner_layering_example1() -> PlanningWorkload {
}),
}
}

/// The shared data workload of #509 with `arrival`.
fn shared_data_workload(arrival: DataArrival) -> DataWorkload {
DataWorkload {
arrival,
data_ingestion_interval: declared(DurationMs(15_000)),
ingestion_volume: Evidence::default(),
ingestion_rate: declared(Rate(1_000_000.0 / 15.0)),
input_cardinality: declared(1_000_000),
distribution: declared(DataDistribution::Zipf),
metric_types: Default::default(),
}
}

/// #509 Example 3, Pattern A: an ad hoc batch of five p99 reports over
/// historical intervals, run once at T (2026-01-01), over mixed data.
fn planner_layering_example3a() -> PlanningWorkload {
const YEAR_MS: u64 = 365 * 24 * 3_600_000;
const T_MS: u64 = 1_767_225_600_000;
let queries = [
("quantile_over_time(0.99, latency_ms[5y])", 5 * YEAR_MS, 0),
("quantile_over_time(0.99, latency_ms[1y])", YEAR_MS, 0),
(
"quantile_over_time(0.99, latency_ms[1y] offset 1y)",
YEAR_MS,
YEAR_MS,
),
(
"quantile_over_time(0.99, latency_ms[1y] offset 2y)",
YEAR_MS,
2 * YEAR_MS,
),
(
"quantile_over_time(0.99, latency_ms[3y] offset 2y)",
3 * YEAR_MS,
2 * YEAR_MS,
),
];
PlanningWorkload {
query_workload: QueryWorkload {
language: QueryLanguage::PromQL,
query_batch: Some(
queries
.into_iter()
.map(|(query, lookback, before_t)| BatchEntry {
query: Query(query.into()),
requirements: QueryRequirements {
accuracy: AccuracyRequirement::Explicit(AccuracyTarget::EpsilonDelta {
epsilon: 0.005,
delta: 0.01,
}),
response_latency: LatencyRequirement::Unspecified,
},
predictability: Predictability::AdHoc,
invocations: 1,
execute_at: Some(TimestampMs(T_MS)),
time_selection: TimeSelection {
scope: QueryTimeScope::Longitudinal,
lookback: Some(DurationMs(lookback)),
as_of: Some(TimestampMs(T_MS - before_t)),
},
})
.collect(),
),
repeating_queries: None,
},
data_workload: Some(shared_data_workload(DataArrival::Mixed)),
}
}

/// #509 Example 3, Pattern B: a p99 panel over the last 5 min, every minute.
fn planner_layering_example3b() -> PlanningWorkload {
PlanningWorkload {
query_workload: QueryWorkload {
language: QueryLanguage::PromQL,
query_batch: None,
repeating_queries: Some(vec![RepeatingEntry {
query: Query("quantile_over_time(0.99, latency_ms[5m])".into()),
demand: RepeatedDemand::FixedInterval(RepetitionInterval(60_000)),
requirements: QueryRequirements {
accuracy: AccuracyRequirement::Explicit(AccuracyTarget::EpsilonDelta {
epsilon: 0.01,
delta: 0.01,
}),
response_latency: LatencyRequirement::ExplicitMaxMs(200.0),
},
predictability: Predictability::Predictable { known_at: None },
time_selection: TimeSelection {
scope: QueryTimeScope::RealTime,
lookback: Some(DurationMs(300_000)),
as_of: None,
},
}]),
},
data_workload: Some(shared_data_workload(DataArrival::ContinuouslyIngesting)),
}
}
54 changes: 48 additions & 6 deletions crates/devtools/tests/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,16 +10,19 @@ use serde_json::Value;

const COMMITTED: &str = "../../tools/dag-viewer/examples/planner-layering-example1.json";

fn generate(extra: &[&str]) -> Value {
/// Every Example 1 candidate (88) is displayed; the default cap is 64.
const EXAMPLE1: [&str; 4] = ["--example", "planner-layering-1", "--max-candidates", "128"];

fn generate(args: &[&str]) -> Value {
let out = std::env::temp_dir().join(format!(
"stage_pipeline_{}_{}.json",
std::process::id(),
extra.len()
args.join("_")
));
let status = Command::new(env!("CARGO_BIN_EXE_stage_pipeline"))
.args(["--example", "planner-layering-1", "--out"])
.args(args)
.arg("--out")
.arg(&out)
.args(extra)
.status()
.unwrap();
assert!(status.success());
Expand All @@ -40,7 +43,7 @@ fn validated_dag(value: &Value, queries: usize) {
/// fixture is current.
#[test]
fn example1_document_is_valid_and_committed_fixture_is_current() {
let document = generate(&[]);
let document = generate(&EXAMPLE1);
let committed: Value = serde_json::from_str(
&std::fs::read_to_string(PathBuf::from(env!("CARGO_MANIFEST_DIR")).join(COMMITTED))
.unwrap(),
Expand Down Expand Up @@ -100,9 +103,48 @@ fn example1_document_is_valid_and_committed_fixture_is_current() {
/// The Cartesian product is cut at `--max-candidates` and says so.
#[test]
fn candidate_cap_is_recorded() {
let document = generate(&["--max-candidates", "5"]);
let document = generate(&["--example", "planner-layering-1", "--max-candidates", "5"]);
let stage1 = &document["stage1_logical_asap"];
assert_eq!(stage1["capped"], true);
assert_eq!(stage1["candidates"].as_array().unwrap().len(), 5);
assert!(stage1["combinations"].as_u64().unwrap() > 5);
}

/// Example 3, Pattern B: the 5-min p99 every minute offers KLL and DDSketch
/// over the whole window and in 1-min tumbling panes, each merged before its
/// estimate.
#[test]
fn example3b_lists_tumbling_candidates() {
let document = generate(&["--example", "planner-layering-3b"]);
let stage1 = &document["stage1_logical_asap"];
let labels: Vec<_> = stage1["candidates"]
.as_array()
.unwrap()
.iter()
.map(|c| c["label"].as_str().unwrap())
.collect();
assert_eq!(
labels,
[
"Q1 exact",
"Q1 Kll",
"Q1 DDSketch",
"Q1 Kll · tumbling 1m panes",
"Q1 DDSketch · tumbling 1m panes"
]
);
for candidate in &stage1["candidates"].as_array().unwrap()[3..] {
let dag: LogicalASAPDAG = serde_json::from_value(candidate["dag"].clone()).unwrap();
let merges = dag
.nodes
.iter()
.filter(|n| {
matches!(
n.payload,
asap_types::ir::export::LogicalASAPOperatorPayload::SummaryMerge
)
})
.count();
assert_eq!(merges, 1, "{}", candidate["label"]);
}
}
Loading