From 9696c8269b804e76daeeb0ac2926c46619363923 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sun, 4 Oct 2026 21:26:59 +0000 Subject: [PATCH 1/2] test(example4): B1 vs B2 as model-relative invariants (Q49) Replace the two ignored Pattern B cost tests, which asserted B1 <= B2 for any workload, with invariants the cost model actually implies: - at 1M series, rebuilding the panes (B2) breaks the 200 ms latency bound; - when the deployment keeps raw data, B2 is cheaper whenever it is valid; - panes larger than the raw samples they cover (15 s sampling, 1-min panes) lose even when the deployment does not keep raw data; - panes smaller than their raw samples (1 s sampling, 10-min panes over 1 h) win when raw data is not kept, and the planner selects B1. run_stages_with plans with given deployment inputs. Co-Authored-By: Claude Opus 5.5 --- .../tests/planner_layering_common/mod.rs | 21 ++-- .../tests/planner_layering_example4.rs | 117 +++++++++++++++--- 2 files changed, 111 insertions(+), 27 deletions(-) diff --git a/crates/integration-tests/tests/planner_layering_common/mod.rs b/crates/integration-tests/tests/planner_layering_common/mod.rs index 56abbab4..73e2ae87 100644 --- a/crates/integration-tests/tests/planner_layering_common/mod.rs +++ b/crates/integration-tests/tests/planner_layering_common/mod.rs @@ -17,20 +17,14 @@ use asap_physical_optimizer::implementation::physical_candidates::PhysicalCandid #[path = "../executor_models/mod.rs"] mod executor_models; -use asap_plan_selection::{plan_stages, Selection}; +use asap_plan_selection::{plan_stages, PlanningModels, Selection}; use asap_types::ir::flat::{flatten, FlatDag}; -use asap_types::ir::physical_export::{PhysicalASAPDAG, PhysicalASAPNodeId}; use asap_types::ir::properties::ExecutionTiming; use asap_types::ir::schema::{FieldDataType, SketchAlgorithm, SketchParams, SketchStatistic}; use asap_types::ir::QueryRoot; use asap_types::ir::{ASAPOp, Operator}; 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, - RepeatedDemand, RepeatingEntry, RepetitionInterval, RootDemand, TimeSelection, TimestampMs, -}; +use asap_types::workload::{AccuracyRequirement, BatchEntry, DataArrival, DataDistribution, DataWorkload, DurationMs, Evidence, EvidenceSource, LatencyRequirement, PlanningWorkload, Predictability, Query, QueryLanguage, QueryRecurrence, QueryRequirements, QueryTimeScope, QueryWorkload, Rate, RepeatedDemand, RepeatingEntry, RepetitionInterval, RootDemand, TimeSelection, TimestampMs}; use executor_models::executor_models; pub type Payload = Operator; @@ -332,6 +326,15 @@ fn export(roots: &[QueryRoot]) -> (FlatDag, Vec) { /// library's `plan_stages`, with each entry's demand (accuracy, recurrence, /// predictability). pub fn run_stages(workload: &PlanningWorkload, roots: Vec) -> Run { + run_stages_with(workload, roots, executor_models()) +} + +/// [`run_stages`] with the deployment inputs `models`. +pub fn run_stages_with( + workload: &PlanningWorkload, + roots: Vec, + models: PlanningModels<'_>, +) -> Run { let (dag, query_roots) = export(&roots); let stage0 = Logical { id: "S0".into(), @@ -348,7 +351,7 @@ pub fn run_stages(workload: &PlanningWorkload, roots: Vec) -> Run { roots.into_iter().enumerate().collect(), &demand, workload.data_workload.as_ref().expect("data workload"), - executor_models(), + models, MAX_CANDIDATES, ) .expect("plans"); diff --git a/crates/integration-tests/tests/planner_layering_example4.rs b/crates/integration-tests/tests/planner_layering_example4.rs index 0f8ce813..16394a55 100644 --- a/crates/integration-tests/tests/planner_layering_example4.rs +++ b/crates/integration-tests/tests/planner_layering_example4.rs @@ -16,9 +16,10 @@ mod planner_layering_common; use asap_types::ir::{ASAPOp, Operator}; use std::collections::{BTreeMap, BTreeSet}; -use asap_types::ir::physical_export::PhysicalASAPNodeId; +use asap_plan_selection::{DeploymentCapabilities, PlanningModels}; +use asap_types::ir::physical_export::{PhysicalASAPNodeId}; use asap_types::ir::schema::SketchAlgorithm; -use asap_types::workload::{DataArrival, PlanningWorkload}; +use asap_types::workload::{DataArrival, DurationMs, PlanningWorkload, Query, Rate, RepeatedDemand, RepetitionInterval}; use planner_layering_common::*; /// Pattern A as given: one ad hoc batch at T over mixed data. @@ -328,26 +329,106 @@ fn stage2_b_sliding_kll_has_two_materialization_options() { assert_eq!(found, BTreeSet::from([IngestionTime, QueryTimeKept])); } -/// B2 rebuilds all five tumbling KLLs at every evaluation, so it costs at least B1 and B3. +// Which of B1 and B2 is cheaper depends on the workload and the deployment, +// not on the materialization alone (Q49). B1 keeps one KLL per series per +// pane for as long as the window; B2 keeps nothing but reads the window's raw +// samples at every evaluation, and pays to keep them when the deployment does +// not. So B1 wins when its panes are smaller than the raw samples they cover +// and those samples would otherwise be kept for the plan. + +/// Pattern B with `series` series sampled every `interval_ms`, asking for the +/// p99 over the last `window_min` minutes every `every_min` minutes. +fn pattern_b_variant( + series: u64, + interval_ms: u64, + window_min: u64, + every_min: u64, +) -> PlanningWorkload { + let mut workload = pattern_b(); + let data = workload.data_workload.as_mut().expect("data workload"); + data.input_cardinality = declared(series); + data.data_ingestion_interval = declared(DurationMs(interval_ms)); + data.ingestion_rate = declared(Rate(series as f64 * 1000.0 / interval_ms as f64)); + let entry = &mut workload + .query_workload + .repeating_queries + .as_mut() + .expect("repeating")[0]; + entry.query = Query(format!( + "quantile_over_time(0.99, latency_ms[{window_min}m])" + )); + entry.demand = + RepeatedDemand::FixedInterval(RepetitionInterval((every_min * MINUTE_MS) as u32)); + entry.time_selection.lookback = Some(DurationMs(window_min * MINUTE_MS)); + workload +} + +/// Plans `workload` on the executor with raw data kept or not by the deployment. +fn run_with_raw_data(workload: &PlanningWorkload, raw_data_retained: bool) -> Run { + let capabilities = DeploymentCapabilities { + raw_data_retained, + ..asap_executor::capabilities() + }; + let models = PlanningModels::builtin().with_capabilities(&capabilities); + run_stages_with(workload, lower_promql(workload), models) +} + +/// The costs of B1 and B2 for the tumbling KLL of `every_min` minutes. +fn b1_b2(run: &Run, every_min: u64) -> (f64, f64) { + let form = Some(WindowForm::Tumbling { + length_ms: every_min * MINUTE_MS, + }); + let options = options_of(run, form, 1); + ( + cost(run, &options[&IngestionTime].0), + cost(run, &options[&NotMaterialized].0), + ) +} + +/// At Pattern B's 1M series, rebuilding the five panes at query time breaks +/// the 200 ms latency bound, so B2 is invalid and B1 is priced. #[test] -#[ignore = "built-in model (#604, Q49): B2 rebuilds the panes in 310 ms, over the 200 ms latency bound, so it is not priced; B1 retains 6 panes of 1M per-series KLLs (6.1 GB) for 768.58 cost/s, and B2 costs 5.17 cost/s, or 45.17 when the deployment does not keep raw data (5 min of raw samples, 320 MB, 40.0 cost/s)"] -fn stage3_b_rebuilding_every_window_costs_most() { +fn stage3_b_rebuilding_a_million_series_breaks_the_latency_bound() { let run = run_promql(&pattern_b()); let options = options_of(&run, tumbling(), 1); - let b2 = cost(&run, &options[&NotMaterialized].0); - for m in [IngestionTime, QueryTimeKept] { - assert!(cost(&run, &options[&m].0) <= b2, "{m:?} vs B2 {b2}"); - } + let b2 = &options[&NotMaterialized].0; + let reason = run.invalid()[b2.as_str()]; + assert!(reason.contains("latency bound"), "{b2}: {reason}"); + assert!(run.cost(&options[&IngestionTime].0).is_some()); } -/// Repeating over arriving data, the built-in models pick B1 among the tumbling options. +/// When the deployment keeps raw data, B2 pays only to read it, while B1 +/// keeps its panes: B2 is cheaper whenever it meets the latency bound. #[test] -#[ignore = "built-in model (#604, Q49): B2 is over the 200 ms latency bound, so it is not priced; B1 costs 768.58 cost/s against B2's 5.17, or 45.17 when the deployment does not keep raw data"] -fn stage3_b_prefers_ingestion_time_tumbling_windows() { - let run = run_promql(&pattern_b()); - let options = options_of(&run, tumbling(), 1); - let b1 = cost(&run, &options[&IngestionTime].0); - for (m, (id, _)) in &options { - assert!(b1 <= cost(&run, id), "B1 {b1} vs {m:?} {}", cost(&run, id)); - } +fn stage3_b_kept_raw_data_makes_rebuilding_cheaper() { + let run = run_with_raw_data(&pattern_b_variant(10_000, 15_000, 5, 1), true); + let (b1, b2) = b1_b2(&run, 1); + assert!(b2 < b1, "B2 {b2} vs B1 {b1}"); +} + +/// Sampled every 15 s, a 1-min pane covers 4 samples per series, fewer bytes +/// than its KLL: B2 is cheaper even when it pays to keep the raw data. +#[test] +fn stage3_b_panes_larger_than_their_raw_data_lose() { + let run = run_with_raw_data(&pattern_b_variant(10_000, 15_000, 5, 1), false); + let (b1, b2) = b1_b2(&run, 1); + assert!(b2 < b1, "B2 {b2} vs B1 {b1}"); +} + +/// Sampled every second, a 10-min pane covers 600 samples per series, far +/// more bytes than its KLL. When the deployment does not keep raw data, B2 +/// must keep the hour's samples, so B1 is cheaper, and the planner selects it. +#[test] +fn stage3_b_panes_smaller_than_their_raw_data_win() { + let run = run_with_raw_data(&pattern_b_variant(1_000, 1_000, 60, 10), false); + let (b1, b2) = b1_b2(&run, 10); + assert!(b1 < b2, "B1 {b1} vs B2 {b2}"); + let options = options_of( + &run, + Some(WindowForm::Tumbling { + length_ms: 10 * MINUTE_MS, + }), + 1, + ); + assert_eq!(run.selection.selected, options[&IngestionTime].0); } From 253be09f79b3c38093c7c06ecafdb3b899cb5d4c Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Thu, 8 Oct 2026 15:00:00 +0000 Subject: [PATCH 2/2] test(example4): read the B1/B2 invariants from Operator payloads Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01W7qG9aFyPij5uWsyAJCxDW --- .../tests/planner_layering_common/mod.rs | 8 +++++++- .../integration-tests/tests/planner_layering_example4.rs | 6 ++++-- 2 files changed, 11 insertions(+), 3 deletions(-) diff --git a/crates/integration-tests/tests/planner_layering_common/mod.rs b/crates/integration-tests/tests/planner_layering_common/mod.rs index 73e2ae87..2be8e2bd 100644 --- a/crates/integration-tests/tests/planner_layering_common/mod.rs +++ b/crates/integration-tests/tests/planner_layering_common/mod.rs @@ -19,12 +19,18 @@ mod executor_models; use asap_plan_selection::{plan_stages, PlanningModels, Selection}; use asap_types::ir::flat::{flatten, FlatDag}; +use asap_types::ir::physical_export::{PhysicalASAPDAG, PhysicalASAPNodeId}; use asap_types::ir::properties::ExecutionTiming; use asap_types::ir::schema::{FieldDataType, SketchAlgorithm, SketchParams, SketchStatistic}; use asap_types::ir::QueryRoot; use asap_types::ir::{ASAPOp, Operator}; 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, RepeatedDemand, RepeatingEntry, RepetitionInterval, RootDemand, TimeSelection, TimestampMs}; +use asap_types::workload::{ + AccuracyRequirement, BatchEntry, DataArrival, DataDistribution, DataWorkload, DurationMs, + Evidence, EvidenceSource, LatencyRequirement, PlanningWorkload, Predictability, Query, + QueryLanguage, QueryRecurrence, QueryRequirements, QueryTimeScope, QueryWorkload, Rate, + RepeatedDemand, RepeatingEntry, RepetitionInterval, RootDemand, TimeSelection, TimestampMs, +}; use executor_models::executor_models; pub type Payload = Operator; diff --git a/crates/integration-tests/tests/planner_layering_example4.rs b/crates/integration-tests/tests/planner_layering_example4.rs index 16394a55..7758c4d1 100644 --- a/crates/integration-tests/tests/planner_layering_example4.rs +++ b/crates/integration-tests/tests/planner_layering_example4.rs @@ -17,9 +17,11 @@ use asap_types::ir::{ASAPOp, Operator}; use std::collections::{BTreeMap, BTreeSet}; use asap_plan_selection::{DeploymentCapabilities, PlanningModels}; -use asap_types::ir::physical_export::{PhysicalASAPNodeId}; +use asap_types::ir::physical_export::PhysicalASAPNodeId; use asap_types::ir::schema::SketchAlgorithm; -use asap_types::workload::{DataArrival, DurationMs, PlanningWorkload, Query, Rate, RepeatedDemand, RepetitionInterval}; +use asap_types::workload::{ + DataArrival, DurationMs, PlanningWorkload, Query, Rate, RepeatedDemand, RepetitionInterval, +}; use planner_layering_common::*; /// Pattern A as given: one ad hoc batch at T over mixed data.