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
105 changes: 74 additions & 31 deletions crates/devtools/src/bin/stage_pipeline.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,8 @@
// cargo run -p asap-devtools --bin stage_pipeline -- \
// --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)
// Patterns A and B; planner-layering-4a: Example 4, Pattern A repeated
// monthly)
// 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 Down Expand Up @@ -58,7 +59,7 @@ use asap_types::workload::{
use serde_json::{json, Value};

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

fn main() {
Expand Down Expand Up @@ -93,6 +94,7 @@ fn run(args: Vec<String>) -> Result<(), String> {
(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("planner-layering-4a"), true) => planner_layering_example4a(),
(Some(other), true) => return Err(format!("unknown example {other}")),
(None, false) => {
let accuracy = match (epsilon, delta) {
Expand Down Expand Up @@ -445,45 +447,52 @@ fn shared_data_workload(arrival: DataArrival) -> DataWorkload {
}
}

const YEAR_MS: u64 = 365 * 24 * 3_600_000;
/// Pattern A's batch time T (2026-01-01T00:00:00Z).
const T_MS: u64 = 1_767_225_600_000;
/// #509 Example 3, Pattern A: five p99 reports, (PromQL, lookback, T − as_of).
const PATTERN_A: [(&str, u64, u64); 5] = [
("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,
),
];

fn pattern_a_requirements() -> QueryRequirements {
QueryRequirements {
accuracy: AccuracyRequirement::Explicit(AccuracyTarget::EpsilonDelta {
epsilon: 0.005,
delta: 0.01,
}),
response_latency: LatencyRequirement::Unspecified,
}
}

/// #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
PATTERN_A
.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,
},
requirements: pattern_a_requirements(),
predictability: Predictability::AdHoc,
invocations: 1,
execute_at: Some(TimestampMs(T_MS)),
Expand All @@ -501,6 +510,40 @@ fn planner_layering_example3a() -> PlanningWorkload {
}
}

/// #509 Example 4, Pattern A repeated monthly and `Predictable { known_at: T }`:
/// each run reads the intervals ending at its own evaluation time, over mixed
/// data. Example 4's other variants are Example 3's workloads (Pattern B is
/// `planner-layering-3b`).
fn planner_layering_example4a() -> PlanningWorkload {
PlanningWorkload {
query_workload: QueryWorkload {
language: QueryLanguage::PromQL,
query_batch: None,
repeating_queries: Some(
PATTERN_A
.into_iter()
.map(|(query, lookback, _)| RepeatingEntry {
query: Query(query.into()),
demand: RepeatedDemand::FixedInterval(RepetitionInterval(
30 * 24 * 3_600_000,
)),
requirements: pattern_a_requirements(),
predictability: Predictability::Predictable {
known_at: Some(TimestampMs(T_MS)),
},
time_selection: TimeSelection {
scope: QueryTimeScope::Longitudinal,
lookback: Some(DurationMs(lookback)),
as_of: None,
},
})
.collect(),
),
},
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 {
Expand Down
41 changes: 41 additions & 0 deletions crates/devtools/tests/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -150,3 +150,44 @@ fn example3b_lists_tumbling_candidates() {
assert_eq!(merges, 1, "{}", candidate["label"]);
}
}

/// Example 4, Pattern A repeated monthly: the same 486 candidates as the ad
/// hoc batch (3a), and none maintained at ingestion time. The windows (1–5 y)
/// are longer than the month between runs and no pane width fits, so
/// nothing is maintainable. The selected plan is 3a's, now amortized over
/// monthly runs instead of one run per hour of horizon.
#[test]
fn example4a_repeats_monthly_with_nothing_maintainable() {
let once = generate(&[
"--example",
"planner-layering-3a",
"--max-candidates",
"600",
]);
let monthly = generate(&[
"--example",
"planner-layering-4a",
"--max-candidates",
"600",
]);
let physical = monthly["stage2_physical_asap"]["candidates"]
.as_array()
.unwrap();
assert_eq!(physical.len(), 486);
assert!(physical
.iter()
.all(|p| !p["label"].as_str().unwrap().contains("ingestion time")));
let selected = |d: &Value| {
d["stage3_selection"]["selected"]
.as_str()
.unwrap()
.to_string()
};
assert_eq!(selected(&monthly), selected(&once));
let cost = |d: &Value| {
d["stage3_selection"]["costs"][selected(d)]["total"]
.as_f64()
.unwrap()
};
assert!(cost(&monthly) < cost(&once));
}
65 changes: 57 additions & 8 deletions crates/integration-tests/tests/planner_layering_common/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -557,7 +557,8 @@ pub fn assert_valid_and_uniquely_named(run: &Run) {
}
}

/// Stage 2 maps the logical candidates one-to-one onto physical ones.
/// Stage 2 gives every logical candidate exactly one all-query-time
/// physical candidate, and possibly more that materialize summaries.
pub fn assert_stage2_bijection(run: &Run) {
let sources: BTreeSet<_> = run
.physical
Expand All @@ -566,7 +567,12 @@ pub fn assert_stage2_bijection(run: &Run) {
.collect();
let logical: BTreeSet<_> = run.logical.iter().map(|c| c.id.as_str()).collect();
assert_eq!(sources, logical);
assert_eq!(run.physical.len(), run.logical.len());
let query_time = run
.physical
.iter()
.filter(|p| p.stage2.materialization.is_empty())
.count();
assert_eq!(query_time, run.logical.len());
}

/// One selected id; every other candidate rejected once, with a reason.
Expand Down Expand Up @@ -699,8 +705,9 @@ pub enum Materialization {
NotMaterialized,
}

/// Pending (needs Stage 2 materialization): the export records only the
/// execution timing, so a query-time node is never kept.
/// Stage 2 runs a node at ingestion time or at query time, recomputed at
/// each evaluation; it does not keep query-time output yet (Example 4 B3),
/// so `QueryTimeKept` does not occur.
pub fn materialization(p: &Physical, node: LogicalASAPNodeId) -> Materialization {
let n = p.dag.nodes.iter().find(|n| n.id == node).expect("node");
match n.output_state.timing {
Expand All @@ -709,10 +716,52 @@ pub fn materialization(p: &Physical, node: LogicalASAPNodeId) -> Materialization
}
}

/// Pending (needs Stage 2 materialization): how long a materialized node's
/// output is kept, in event time. `None` until Stage 2 records retention.
pub fn retention_ms(_p: &Physical, _node: LogicalASAPNodeId) -> Option<u64> {
None
/// How long a materialized node's output is kept, in event time. The export
/// records no retention, so this derives it as Stage 3 prices it
/// (`stage3-cost-model.md`): ingestion-time work read at query time through
/// a merge of `N` panes of width `w` keeps `(N + 1) · w`; read directly, the
/// window being built and the completed one, `2 · window`. Taken over every
/// query-time reader the node's ingestion-time work feeds. `None` for a
/// query-time node.
pub fn retention_ms(p: &Physical, node: LogicalASAPNodeId) -> Option<u64> {
if !runs_at_ingestion(p, node) {
return None;
}
// The longest raw range an ingestion-time node reads.
let window = |id: LogicalASAPNodeId| {
closure(&p.dag, id)
.into_iter()
.filter_map(|n| match p.dag.payload(n) {
LogicalASAPOperatorPayload::Relational {
operator: NonASAPOpKind::TimeRange { range, .. },
} => Some(range.as_millis() as u64),
_ => None,
})
.max()
.unwrap_or(0)
};
let mut kept = 0;
let mut stack = vec![node];
let mut seen = HashSet::new();
while let Some(id) = stack.pop() {
if !seen.insert(id) {
continue;
}
for consumer in p.dag.consumers(id) {
if runs_at_ingestion(p, consumer) {
stack.push(consumer);
} else if matches!(
p.dag.payload(consumer),
LogicalASAPOperatorPayload::SummaryMerge
) {
let panes = p.dag.producers(consumer).len() as u64;
kept = kept.max((panes + 1) * window(id));
} else {
kept = kept.max(2 * window(id));
}
}
}
Some(kept)
}

pub fn runs_at_ingestion(p: &Physical, node: LogicalASAPNodeId) -> bool {
Expand Down
Loading