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
6 changes: 4 additions & 2 deletions crates/devtools/src/bin/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,8 +42,8 @@ use asap_types::ir::{OperatorNode, QueryRoot};
use asap_types::types::AccuracyTarget;
use asap_types::workload::{
AccuracyRequirement, BatchEntry, DataArrival, DataDistribution, DataWorkload, DurationMs,
Evidence, EvidenceSource, LatencyRequirement, PlanningWorkload, Predictability, Query,
QueryLanguage, QueryRecurrence, QueryRequirements, QueryTimeScope, QueryWorkload, Rate,
Evidence, EvidenceSource, LatencyRequirement, MetricType, PlanningWorkload, Predictability,
Query, QueryLanguage, QueryRecurrence, QueryRequirements, QueryTimeScope, QueryWorkload, Rate,
RepeatedDemand, RepeatingEntry, RepetitionInterval, RootDemand, TimeSelection,
};
use serde_json::{json, Value};
Expand Down Expand Up @@ -406,6 +406,8 @@ fn planner_layering_example1() -> PlanningWorkload {
ingestion_rate: declared(Rate(1_000_000.0 / 15.0)),
input_cardinality: declared(1_000_000),
distribution: declared(DataDistribution::Zipf),
// `http_requests_total` is a counter: its samples are never negative.
metric_types: [("http_requests_total".into(), MetricType::Counter)].into(),
}),
}
}
80 changes: 61 additions & 19 deletions crates/integration-tests/tests/planner_layering_example1.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,9 +24,9 @@ use asap_types::ir::{ASAPOp, Operator};
use asap_types::types::AccuracyTarget;
use asap_types::workload::{
AccuracyRequirement, DataArrival, DataDistribution, DataWorkload, DurationMs, Evidence,
EvidenceSource, LatencyRequirement, PlanningWorkload, Predictability, Query, QueryLanguage,
QueryRequirements, QueryTimeScope, QueryWorkload, Rate, RepeatedDemand, RepeatingEntry,
RepetitionInterval, TimeSelection,
EvidenceSource, LatencyRequirement, MetricType, PlanningWorkload, Predictability, Query,
QueryLanguage, QueryRequirements, QueryTimeScope, QueryWorkload, Rate, RepeatedDemand,
RepeatingEntry, RepetitionInterval, TimeSelection,
};

type Payload = Operator<NodeId>;
Expand Down Expand Up @@ -337,6 +337,8 @@ fn example1_workload() -> PlanningWorkload {
ingestion_rate: declared(Rate(1_000_000.0 / 15.0)),
input_cardinality: declared(1_000_000),
distribution: declared(DataDistribution::Zipf),
// `http_requests_total` is a counter: its samples are never negative.
metric_types: [("http_requests_total".into(), MetricType::Counter)].into(),
}),
}
}
Expand Down Expand Up @@ -515,12 +517,31 @@ fn snake_case(name: &str) -> String {
out
}

/// Example 1 with `http_requests_total` declared `metric_type`, or undeclared.
fn example1_workload_with(metric_type: Option<MetricType>) -> PlanningWorkload {
let mut workload = example1_workload();
let data = workload.data_workload.as_mut().expect("data workload");
data.metric_types = metric_type
.map(|t| [("http_requests_total".to_string(), t)].into())
.unwrap_or_default();
workload
}

fn pipeline() -> (
PlanningWorkload,
Vec<LogicalCandidate>,
Vec<PhysicalCandidate>,
) {
let workload = example1_workload();
pipeline_for(example1_workload())
}

fn pipeline_for(
workload: PlanningWorkload,
) -> (
PlanningWorkload,
Vec<LogicalCandidate>,
Vec<PhysicalCandidate>,
) {
let logical = stage1_logical_asap(&workload, &stage0_logical(&workload));
let physical = stage2_physical(&workload, &logical);
(workload, logical, physical)
Expand Down Expand Up @@ -898,21 +919,18 @@ fn compile_in_runtime(p: &PhysicalCandidate) -> Result<(), String> {
.map_err(|e| format!("{} ({}): {e}", p.id, p.label))
}

/// Runtime capability check (added by the implementer, not part of the
/// spec): the physical planner compiles every candidate Stage 3 finds valid,
/// and rejects the invalid ones (Count-Min over weights not proven
/// non-negative) for the same reason Stage 3 gives.
#[test]
fn stage2_runtime_compiles_exactly_the_candidates_stage3_finds_valid() {
let (workload, _, physical) = pipeline();
/// The physical planner compiles every candidate Stage 3 finds valid, and
/// rejects the invalid ones (Count-Min over weights not proven non-negative)
/// for the same reason Stage 3 gives. Returns the number of invalid ones.
fn assert_runtime_agrees_with_stage3(workload: PlanningWorkload) -> usize {
let (workload, _, physical) = pipeline_for(workload);
let selection = stage3_select(&workload, &physical, PlanningModels::builtin());
let invalid: BTreeMap<_, _> = selection
.rejected
.iter()
.filter(|r| !r.valid)
.map(|r| (r.id.as_str(), r.reason.as_str()))
.collect();
assert_eq!(invalid.len(), 24, "the Count-Min + heap candidates");
for p in &physical {
let compiled = compile_in_runtime(p);
match invalid.get(p.id.as_str()) {
Expand All @@ -931,6 +949,28 @@ fn stage2_runtime_compiles_exactly_the_candidates_stage3_finds_valid() {
}
}
}
invalid.len()
}

/// Runtime capability check (added by the implementer, not part of the
/// spec): with `http_requests_total` declared a counter, every candidate,
/// Count-Min + heap included, is valid in Stage 3 and compiles.
#[test]
fn stage2_runtime_compiles_every_candidate_over_a_declared_counter() {
assert_eq!(assert_runtime_agrees_with_stage3(example1_workload()), 0);
}

/// Without the counter declaration, or with a gauge, the Count-Min + heap
/// candidates are invalid in Stage 3 and the runtime rejects them.
#[test]
fn stage2_count_min_needs_a_counter_declaration() {
for metric_type in [None, Some(MetricType::Gauge)] {
assert_eq!(
assert_runtime_agrees_with_stage3(example1_workload_with(metric_type)),
24,
"{metric_type:?}: the Count-Min + heap candidates"
);
}
}

/// A CountSketch+heap top-k readout compiles: the IR's derived readout schema
Expand Down Expand Up @@ -1008,7 +1048,7 @@ fn stage3_selects_cheapest_valid() {

/// Per-second cost keeps Example 1's ranking: both panels repeat every
/// 10 s and everything runs at query time, so every candidate costs 0.1 ×
/// its per-evaluation cost, and P58 still wins at 52.201 × 0.1 per second.
/// its per-evaluation cost, and P60 still wins at 46.201 × 0.1 per second.
#[test]
fn stage3_per_second_cost_keeps_the_ranking() {
let (workload, _, physical) = pipeline();
Expand All @@ -1024,14 +1064,14 @@ fn stage3_per_second_cost_keeps_the_ranking() {
entry.demand = RepeatedDemand::FixedInterval(RepetitionInterval(1_000));
}
let per_evaluation = stage3_select(&every_second, &physical, PlanningModels::builtin());
assert_eq!(per_second.selected, "P58");
assert_eq!(per_evaluation.selected, "P58");
assert_eq!(per_second.selected, "P60");
assert_eq!(per_evaluation.selected, "P60");
for (id, cost) in &per_second.costs {
let expected = 0.1 * per_evaluation.costs[id].total;
assert!((cost.total - expected).abs() <= 1e-9 * expected, "{id}");
}
let best = per_second.costs["P58"].total;
assert!((best - 5.2201).abs() < 1e-3, "{best}");
let best = per_second.costs["P60"].total;
assert!((best - 4.6201).abs() < 1e-3, "{best}");
}

/// Every node is charged exactly once, so a shared input is costed once for
Expand Down Expand Up @@ -1069,8 +1109,8 @@ fn stage3_charges_each_node_once() {
}

/// Sharing the input never costs more than reading it separately, for the
/// same local choices. Count-Min + heap is invalid here and has no cost
/// (Hydra: see `stage1_q2_summary_families_are_heap_sketches_and_hydra`).
/// same local choices (Hydra: see
/// `stage1_q2_summary_families_are_heap_sketches_and_hydra`).
#[test]
fn stage3_shared_input_is_not_costlier() {
let (workload, _, physical) = pipeline();
Expand All @@ -1085,7 +1125,9 @@ fn stage3_shared_input_is_not_costlier() {
.collect();
for option in [
Q2Option::Exact,
Q2Option::CountMinHeapPerJob,
Q2Option::CountSketchHeapPerJob,
Q2Option::WholeCountMinHeapPerJob,
Q2Option::WholeCountSketchHeapPerJob,
] {
assert!(
Expand Down
Loading