From 4ab777a88256f9a463a6903e63560c6f10e7600e Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 7 Oct 2026 03:21:13 +0000 Subject: [PATCH] test(e2e): answer #509's PromQL examples like the Planner executor For Examples 1, 3a, 3b, 4a and 4b: plan the workload, start the data_plane binary against a Prometheus holding generated samples, Remote Write the same samples, and require every planned query to be answered by the plan with the Planner executor's result for the workload planned without ingestion-time state. Exact families must match exactly; Example 3a/4a's five-year KLL windows exceed k and KLL compaction is randomized, so they must agree within twice the 0.005 rank target. Co-Authored-By: Claude Opus 5.5 --- README.md | 8 + data_plane/tests/e2e_examples.rs | 411 +++++++++++++++++++++++ examples/planner-layering-example3a.json | 152 +++++++++ examples/planner-layering-example4a.json | 177 ++++++++++ 4 files changed, 748 insertions(+) create mode 100644 data_plane/tests/e2e_examples.rs create mode 100644 examples/planner-layering-example3a.json create mode 100644 examples/planner-layering-example4a.json diff --git a/README.md b/README.md index c42eb009..ca492145 100644 --- a/README.md +++ b/README.md @@ -61,6 +61,14 @@ cargo run -p data_plane -- --plan plan.json --prometheus-url http://localhost:90 `--capabilities caps.json` replaces the default capabilities, which are the Planner executor's capabilities without state kept across query evaluations. +## Examples + +`examples/` holds #509's PromQL example workloads (1, 3a, 3b, 4a, 4b) as +Planner `PlanningWorkload`s. `data_plane/tests/e2e_examples.rs` plans each +one, runs the `data_plane` binary against a Prometheus holding generated +samples, Remote Writes the same samples, and checks every planned query +against the Planner executor running the workload entirely at query time. + ## Development ```sh diff --git a/data_plane/tests/e2e_examples.rs b/data_plane/tests/e2e_examples.rs new file mode 100644 index 00000000..3fefc22e --- /dev/null +++ b/data_plane/tests/e2e_examples.rs @@ -0,0 +1,411 @@ +//! End to end over the #509 PromQL examples: plan each workload, run the +//! `data_plane` binary against a Prometheus holding the samples, Remote +//! Write the same samples to it, and check that every planned query is +//! answered by the plan with what the Planner executor computes when the whole +//! workload runs at query time over those samples. + +mod common; + +use std::collections::BTreeMap; +use std::path::{Path, PathBuf}; +use std::process::{Child, Command}; +use std::sync::Arc; +use std::time::{Duration, Instant}; + +use asap_executor::physical_planner::frontier_from_timing; +use axum::extract::Query; +use axum::routing::get; +use axum::{Json, Router}; +use control_plane::{backend_capabilities, plan_workload, InstalledPlan}; +use data_plane::precompute::RawSample; +use data_plane::query::{QueryEngine, ResultSample}; +use data_plane::remote_write::{encode, Label, Sample, TimeSeries, WriteRequest}; +use data_plane::sds::{Sds, SummaryStore}; +use planner_types::deployment::DeploymentCapabilities; +use planner_types::workload::PlanningWorkload; +use serde_json::{json, Value}; + +const MINUTE: i64 = 60_000; +const DAY: i64 = 24 * 60 * MINUTE; +/// #509 Example 3/4 Pattern A's batch time T (2026-01-01T00:00:00Z). +const PATTERN_A_T: i64 = 1_767_225_600_000; + +fn workload(example: &str) -> PlanningWorkload { + let path = Path::new(env!("CARGO_MANIFEST_DIR")) + .join(format!("../examples/planner-layering-{example}.json")); + serde_json::from_str(&std::fs::read_to_string(path).unwrap()).unwrap() +} + +/// `samples` of each series, every `step` ms in `[from, to)`; series `i`'s +/// value at `t` is `f(i, t)`. +fn samples( + series: &[BTreeMap], + from: i64, + to: i64, + step: i64, + f: impl Fn(usize, i64) -> f64, +) -> Vec { + let series: Vec<_> = series.iter().cloned().map(Arc::new).collect(); + (from..to) + .step_by(step as usize) + .flat_map(|t| { + series + .iter() + .enumerate() + .map(move |(i, labels)| (i, labels.clone(), t)) + }) + .map(|(i, labels, timestamp_ms)| RawSample { + labels, + timestamp_ms, + value: f(i, timestamp_ms), + }) + .collect() +} + +fn series(metric: &str, labels: &[(&str, &str)]) -> BTreeMap { + let mut map = BTreeMap::from([("__name__".to_string(), metric.to_string())]); + map.extend(labels.iter().map(|(k, v)| (k.to_string(), v.to_string()))); + map +} + +/// The answer of the Planner executor for `query` at `time_ms` when the +/// workload is planned without ingestion-time state: every input is a raw +/// range over `all`. +fn reference( + workload: &PlanningWorkload, + query: &str, + time_ms: i64, + all: &[RawSample], +) -> Vec { + let capabilities = DeploymentCapabilities { + ingestion_time: false, + ..backend_capabilities() + }; + let plan = plan_workload(workload, &capabilities).unwrap(); + let sds = Sds::from_dag(&plan.dag.dag, &frontier_from_timing(&plan.dag.dag).unwrap()).unwrap(); + assert!(sds.definitions.is_empty()); + let engine = QueryEngine::new(&plan, &sds); + let store = SummaryStore::new(); + let (root, ranges) = engine.raw_ranges(query, time_ms, &store).unwrap(); + let raw = ranges + .iter() + .map(|range| { + let samples = all + .iter() + .filter(|s| { + s.labels["__name__"] == range.metric + && range.start_ms < s.timestamp_ms + && s.timestamp_ms <= range.end_ms + }) + .cloned() + .collect(); + (range.node, samples) + }) + .collect(); + sorted(engine.evaluate(root, time_ms, &store, &raw).unwrap()) +} + +fn sorted(mut samples: Vec) -> Vec { + samples.sort_by(|a, b| a.labels.cmp(&b.labels)); + samples +} + +/// A Prometheus holding `all`: range selectors `{__name__="m"}[Nms]` at +/// `time` return the samples in `[time - N, time]`, as Prometheus 2 does. +async fn prometheus(all: Arc>) -> String { + let handler = move |Query(params): Query>| { + let all = all.clone(); + async move { + let query = ¶ms["query"]; + let (metric, rest) = query + .strip_prefix("{__name__=\"") + .and_then(|rest| rest.split_once("\"}[")) + .expect("the data plane only asks for raw ranges"); + let range: i64 = rest.strip_suffix("ms]").unwrap().parse().unwrap(); + let time = (params["time"].parse::().unwrap() * 1000.0).round() as i64; + let mut by_series: BTreeMap<&BTreeMap, Vec> = BTreeMap::new(); + for s in all.iter() { + if s.labels["__name__"] == metric + && time - range <= s.timestamp_ms + && s.timestamp_ms <= time + { + by_series + .entry(&s.labels) + .or_default() + .push(json!([s.timestamp_ms as f64 / 1000.0, s.value.to_string()])); + } + } + let result: Vec = by_series + .into_iter() + .map(|(labels, values)| json!({"metric": labels, "values": values})) + .collect(); + Json(json!({"status": "success", "data": {"resultType": "matrix", "result": result}})) + } + }; + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("http://{}", listener.local_addr().unwrap()); + let app = Router::new().route("/api/v1/query", get(handler)); + tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + url +} + +/// The `data_plane` binary serving `plan`, killed on drop. +struct DataPlaneProcess { + child: Child, + url: String, + _dir: TempDir, +} + +impl Drop for DataPlaneProcess { + fn drop(&mut self) { + let _ = self.child.kill(); + let _ = self.child.wait(); + } +} + +struct TempDir(PathBuf); + +impl Drop for TempDir { + fn drop(&mut self) { + let _ = std::fs::remove_dir_all(&self.0); + } +} + +async fn start(plan: &InstalledPlan, prometheus_url: &str) -> DataPlaneProcess { + let dir = std::env::temp_dir().join(format!( + "data-plane-e2e-{}-{}", + std::process::id(), + rand_suffix() + )); + std::fs::create_dir_all(&dir).unwrap(); + let path = dir.join("plan.json"); + plan.save(&path).unwrap(); + let port = std::net::TcpListener::bind("127.0.0.1:0") + .unwrap() + .local_addr() + .unwrap() + .port(); + let child = Command::new(env!("CARGO_BIN_EXE_data_plane")) + .args(["--plan", path.to_str().unwrap()]) + .args(["--prometheus-url", prometheus_url]) + .args(["--listen", &format!("127.0.0.1:{port}")]) + .args(["--allowed-lateness-ms", "0"]) + .spawn() + .unwrap(); + let process = DataPlaneProcess { + child, + url: format!("http://127.0.0.1:{port}"), + _dir: TempDir(dir), + }; + let deadline = Instant::now() + Duration::from_secs(30); + let client = reqwest::Client::new(); + while client + .get(format!("{}/-/ready", process.url)) + .send() + .await + .is_err() + { + assert!(Instant::now() < deadline, "data plane did not start"); + tokio::time::sleep(Duration::from_millis(50)).await; + } + process +} + +fn rand_suffix() -> u128 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_nanos() +} + +async fn remote_write(url: &str, all: &[RawSample]) { + let mut by_series: BTreeMap<&BTreeMap, Vec> = BTreeMap::new(); + for s in all { + by_series.entry(&s.labels).or_default().push(Sample { + value: s.value, + timestamp: s.timestamp_ms, + }); + } + let timeseries = by_series + .into_iter() + .map(|(labels, samples)| TimeSeries { + labels: labels + .iter() + .map(|(name, value)| Label { + name: name.clone(), + value: value.clone(), + }) + .collect(), + samples, + }) + .collect(); + let response = reqwest::Client::new() + .post(format!("{url}/api/v1/write")) + .body(encode(&WriteRequest { timeseries })) + .send() + .await + .unwrap(); + assert_eq!(response.status(), 204); +} + +/// The plan's answer to `query` at `time_ms`, retried until the plan answers +/// (stored panes are built asynchronously after Remote Write). +async fn planned_answer(url: &str, query: &str, time_ms: i64) -> Vec { + let client = reqwest::Client::new(); + let deadline = Instant::now() + Duration::from_secs(30); + loop { + let response = client + .get(format!("{url}/api/v1/query")) + .query(&[ + ("query", query), + ("time", &(time_ms as f64 / 1000.0).to_string()), + ]) + .send() + .await + .unwrap(); + let execution = response.headers().get("x-asap-execution").cloned(); + let reason = response.headers().get("x-asap-forward-reason").cloned(); + let body: Value = response.json().await.unwrap(); + if execution.as_ref().is_some_and(|v| v == "plan") { + let result = body["data"]["result"].as_array().unwrap(); + return sorted( + result + .iter() + .map(|s| ResultSample { + labels: serde_json::from_value(s["metric"].clone()).unwrap(), + value: s["value"][1].as_str().unwrap().parse().unwrap(), + }) + .collect(), + ); + } + assert!( + Instant::now() < deadline, + "{query} at {time_ms} was not answered by the plan: {reason:?}" + ); + tokio::time::sleep(Duration::from_millis(50)).await; + } +} + +/// `actual` and `expected` hold the same series with values within +/// `tolerance` (plus float rounding). +fn assert_same(query: &str, actual: &[ResultSample], expected: &[ResultSample], tolerance: f64) { + assert!(!expected.is_empty(), "{query}: the reference has no result"); + assert_eq!( + actual.len(), + expected.len(), + "{query}: {actual:?} vs {expected:?}" + ); + for (a, e) in actual.iter().zip(expected) { + assert_eq!(a.labels, e.labels, "{query}"); + let tolerance = tolerance + 1e-9 * e.value.abs().max(1.0); + assert!( + (a.value - e.value).abs() <= tolerance, + "{query}: {a:?} vs {e:?}" + ); + } +} + +/// Runs every planned query of `example` at `time_ms` and compares it with +/// the reference, within `tolerance`. Returns how many stored summary +/// definitions the selected plan has. +async fn check_example( + example: &str, + capabilities: DeploymentCapabilities, + all: Vec, + time_ms: i64, + tolerance: f64, +) -> usize { + let workload = workload(example); + let plan = plan_workload(&workload, &capabilities).unwrap(); + let all = Arc::new(all); + let process = start(&plan, &prometheus(all.clone()).await).await; + remote_write(&process.url, &all).await; + for query in &plan.queries { + let expected = reference(&workload, query, time_ms, &all); + let actual = planned_answer(&process.url, query, time_ms).await; + assert_same(query, &actual, &expected, tolerance); + } + let dag = &plan.dag.dag; + let sds = Sds::from_dag(dag, &frontier_from_timing(dag).unwrap()).unwrap(); + sds.definitions.len() +} + +/// Example 1: per-job request rates exactly and top-k per job with a +/// Count-Min sketch, over three counters. +#[tokio::test] +async fn example1() { + let series = [ + series("http_requests_total", &[("job", "api"), ("instance", "1")]), + series("http_requests_total", &[("job", "api"), ("instance", "2")]), + series("http_requests_total", &[("job", "web"), ("instance", "1")]), + ]; + let all = samples(&series, 0, 10 * MINUTE, 15_000, |i, t| { + (i as f64 + 1.0) * t as f64 / 1000.0 + }); + check_example("example1", backend_capabilities(), all, 5 * MINUTE, 0.0).await; +} + +/// Example 3a: five p99 reports over historical years at T, merged from +/// one-year KLL panes at query time; one sample per series per day. +/// The five-year windows hold more samples than the KLL's k, and its +/// compaction is randomized, so two runs agree only within the target: each +/// is within rank error 0.005 of the exact p99, about 5 for values spread +/// evenly over 0..997. +#[tokio::test] +async fn example3a() { + let series = [ + series("latency_ms", &[("service", "a")]), + series("latency_ms", &[("service", "b")]), + ]; + let all = samples( + &series, + PATTERN_A_T - 6 * 365 * DAY, + PATTERN_A_T, + DAY, + |i, t| (i * 10_000) as f64 + ((t / DAY) % 997) as f64, + ); + check_example("example3a", backend_capabilities(), all, PATTERN_A_T, 10.0).await; +} + +/// Example 4a: Example 3a's reports repeated monthly, evaluated at T. +#[tokio::test] +async fn example4a() { + let series = [series("latency_ms", &[("service", "a")])]; + let all = samples( + &series, + PATTERN_A_T - 6 * 365 * DAY, + PATTERN_A_T, + DAY, + |_, t| ((t / DAY) % 997) as f64, + ); + check_example("example4a", backend_capabilities(), all, PATTERN_A_T, 10.0).await; +} + +/// Example 3b: a p99 panel over the last 5 minutes; KLL is exact below k. +#[tokio::test] +async fn example3b() { + let series = [series("latency_ms", &[("service", "a")])]; + let all = samples(&series, 0, 20 * MINUTE, 15_000, |_, t| { + ((t / 15_000) % 37) as f64 + }); + check_example("example3b", backend_capabilities(), all, 10 * MINUTE, 0.0).await; +} + +/// Example 4b: the hourly p99 on a deployment that keeps no raw data, from +/// six ten-minute KLL panes built from Remote Write. +#[tokio::test] +async fn example4b() { + let series = [ + series("latency_ms", &[("service", "a")]), + series("latency_ms", &[("service", "b")]), + ]; + let all = samples(&series, MINUTE, 75 * MINUTE, 15_000, |i, t| { + (i * 1000) as f64 + ((t / 15_000) % 240) as f64 + }); + let capabilities = DeploymentCapabilities { + raw_data_retained: false, + ..backend_capabilities() + }; + let definitions = check_example("example4b", capabilities, all, 70 * MINUTE, 0.0).await; + assert_eq!(definitions, 1); +} diff --git a/examples/planner-layering-example3a.json b/examples/planner-layering-example3a.json new file mode 100644 index 00000000..1229b0ce --- /dev/null +++ b/examples/planner-layering-example3a.json @@ -0,0 +1,152 @@ +{ + "query_workload": { + "language": "promql", + "query_batch": [ + { + "query": "quantile_over_time(0.99, latency_ms[5y])", + "requirements": { + "accuracy": { + "explicit": { + "EpsilonDelta": { + "epsilon": 0.005, + "delta": 0.01 + } + } + }, + "response_latency": "unspecified" + }, + "predictability": "ad_hoc", + "invocations": 1, + "execute_at": 1767225600000, + "time_selection": { + "scope": "longitudinal", + "lookback": 157680000000, + "as_of": 1767225600000 + } + }, + { + "query": "quantile_over_time(0.99, latency_ms[1y])", + "requirements": { + "accuracy": { + "explicit": { + "EpsilonDelta": { + "epsilon": 0.005, + "delta": 0.01 + } + } + }, + "response_latency": "unspecified" + }, + "predictability": "ad_hoc", + "invocations": 1, + "execute_at": 1767225600000, + "time_selection": { + "scope": "longitudinal", + "lookback": 31536000000, + "as_of": 1767225600000 + } + }, + { + "query": "quantile_over_time(0.99, latency_ms[1y] offset 1y)", + "requirements": { + "accuracy": { + "explicit": { + "EpsilonDelta": { + "epsilon": 0.005, + "delta": 0.01 + } + } + }, + "response_latency": "unspecified" + }, + "predictability": "ad_hoc", + "invocations": 1, + "execute_at": 1767225600000, + "time_selection": { + "scope": "longitudinal", + "lookback": 31536000000, + "as_of": 1735689600000 + } + }, + { + "query": "quantile_over_time(0.99, latency_ms[1y] offset 2y)", + "requirements": { + "accuracy": { + "explicit": { + "EpsilonDelta": { + "epsilon": 0.005, + "delta": 0.01 + } + } + }, + "response_latency": "unspecified" + }, + "predictability": "ad_hoc", + "invocations": 1, + "execute_at": 1767225600000, + "time_selection": { + "scope": "longitudinal", + "lookback": 31536000000, + "as_of": 1704153600000 + } + }, + { + "query": "quantile_over_time(0.99, latency_ms[3y] offset 2y)", + "requirements": { + "accuracy": { + "explicit": { + "EpsilonDelta": { + "epsilon": 0.005, + "delta": 0.01 + } + } + }, + "response_latency": "unspecified" + }, + "predictability": "ad_hoc", + "invocations": 1, + "execute_at": 1767225600000, + "time_selection": { + "scope": "longitudinal", + "lookback": 94608000000, + "as_of": 1704153600000 + } + } + ], + "repeating_queries": null + }, + "data_workload": { + "arrival": "mixed", + "data_ingestion_interval": { + "value": 15000, + "source": "declared", + "observed_at_ms": null, + "valid_for_ms": null + }, + "ingestion_volume": { + "value": null, + "source": "unknown", + "observed_at_ms": null, + "valid_for_ms": null + }, + "ingestion_rate": { + "value": 66666.66666666667, + "source": "declared", + "observed_at_ms": null, + "valid_for_ms": null + }, + "input_cardinality": { + "value": 1000000, + "source": "declared", + "observed_at_ms": null, + "valid_for_ms": null + }, + "distribution": { + "value": "zipf", + "source": "declared", + "observed_at_ms": null, + "valid_for_ms": null + }, + "metric_types": {} + } +} diff --git a/examples/planner-layering-example4a.json b/examples/planner-layering-example4a.json new file mode 100644 index 00000000..ac7db6b0 --- /dev/null +++ b/examples/planner-layering-example4a.json @@ -0,0 +1,177 @@ +{ + "query_workload": { + "language": "promql", + "query_batch": null, + "repeating_queries": [ + { + "query": "quantile_over_time(0.99, latency_ms[5y])", + "demand": { + "fixed_interval": 2592000000 + }, + "requirements": { + "accuracy": { + "explicit": { + "EpsilonDelta": { + "epsilon": 0.005, + "delta": 0.01 + } + } + }, + "response_latency": "unspecified" + }, + "predictability": { + "predictable": { + "known_at": 1767225600000 + } + }, + "time_selection": { + "scope": "longitudinal", + "lookback": 157680000000, + "as_of": null + } + }, + { + "query": "quantile_over_time(0.99, latency_ms[1y])", + "demand": { + "fixed_interval": 2592000000 + }, + "requirements": { + "accuracy": { + "explicit": { + "EpsilonDelta": { + "epsilon": 0.005, + "delta": 0.01 + } + } + }, + "response_latency": "unspecified" + }, + "predictability": { + "predictable": { + "known_at": 1767225600000 + } + }, + "time_selection": { + "scope": "longitudinal", + "lookback": 31536000000, + "as_of": null + } + }, + { + "query": "quantile_over_time(0.99, latency_ms[1y] offset 1y)", + "demand": { + "fixed_interval": 2592000000 + }, + "requirements": { + "accuracy": { + "explicit": { + "EpsilonDelta": { + "epsilon": 0.005, + "delta": 0.01 + } + } + }, + "response_latency": "unspecified" + }, + "predictability": { + "predictable": { + "known_at": 1767225600000 + } + }, + "time_selection": { + "scope": "longitudinal", + "lookback": 31536000000, + "as_of": null + } + }, + { + "query": "quantile_over_time(0.99, latency_ms[1y] offset 2y)", + "demand": { + "fixed_interval": 2592000000 + }, + "requirements": { + "accuracy": { + "explicit": { + "EpsilonDelta": { + "epsilon": 0.005, + "delta": 0.01 + } + } + }, + "response_latency": "unspecified" + }, + "predictability": { + "predictable": { + "known_at": 1767225600000 + } + }, + "time_selection": { + "scope": "longitudinal", + "lookback": 31536000000, + "as_of": null + } + }, + { + "query": "quantile_over_time(0.99, latency_ms[3y] offset 2y)", + "demand": { + "fixed_interval": 2592000000 + }, + "requirements": { + "accuracy": { + "explicit": { + "EpsilonDelta": { + "epsilon": 0.005, + "delta": 0.01 + } + } + }, + "response_latency": "unspecified" + }, + "predictability": { + "predictable": { + "known_at": 1767225600000 + } + }, + "time_selection": { + "scope": "longitudinal", + "lookback": 94608000000, + "as_of": null + } + } + ] + }, + "data_workload": { + "arrival": "mixed", + "data_ingestion_interval": { + "value": 15000, + "source": "declared", + "observed_at_ms": null, + "valid_for_ms": null + }, + "ingestion_volume": { + "value": null, + "source": "unknown", + "observed_at_ms": null, + "valid_for_ms": null + }, + "ingestion_rate": { + "value": 66666.66666666667, + "source": "declared", + "observed_at_ms": null, + "valid_for_ms": null + }, + "input_cardinality": { + "value": 1000000, + "source": "declared", + "observed_at_ms": null, + "valid_for_ms": null + }, + "distribution": { + "value": "zipf", + "source": "declared", + "observed_at_ms": null, + "valid_for_ms": null + }, + "metric_types": {} + } +}