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": {} + } +}