diff --git a/.github/workflows/promql-compliance.yml b/.github/workflows/promql-compliance.yml new file mode 100644 index 00000000..922cabb3 --- /dev/null +++ b/.github/workflows/promql-compliance.yml @@ -0,0 +1,37 @@ +name: PromQL compliance + +on: + workflow_dispatch: + +permissions: + contents: read + +jobs: + differential: + runs-on: ubuntu-latest + timeout-minutes: 90 + steps: + - uses: actions/checkout@v4 + with: + path: ASAPQuery-backend + - uses: actions/checkout@v4 + with: + repository: ProjectASAP/ASAPCollector + path: ASAPCollector + - uses: actions/checkout@v4 + with: + repository: ProjectASAP/asap_sketchlib + path: asap_sketchlib + - uses: actions/setup-go@v5 + with: + go-version: "1.25.8" + - name: Run every differential corpus + working-directory: ASAPQuery-backend/promql-compliance/runner + run: make run-all REPORT_DIR="$GITHUB_WORKSPACE/artifacts/reports" LOGS_DIR="$GITHUB_WORKSPACE/artifacts/logs" + - name: Upload reports and service logs + if: always() + uses: actions/upload-artifact@v4 + with: + name: promql-compliance-evidence + path: artifacts + if-no-files-found: warn diff --git a/control_plane/Dockerfile b/control_plane/Dockerfile index 3bf61346..936f07c5 100644 --- a/control_plane/Dockerfile +++ b/control_plane/Dockerfile @@ -15,6 +15,7 @@ # --build-context asap-gorilla-rust=/path/ASAPCollector/asap-gorilla-rust \ # -t asap/control-plane:dev . FROM rust:1.90-bookworm AS build +ARG LOCAL_DEP_OVERRIDES=false WORKDIR /src # protoc for control_plane/build.rs (prost-build/tonic-build compile @@ -27,11 +28,26 @@ RUN apt-get update && \ # Same sibling-checkout layout as data_plane/Dockerfile so workspace path-deps # resolve during the build. COPY . ASAPQuery-backend +# hadolint ignore=DL3022 COPY --from=asap-precompute-rs . ASAPCollector/asap-precompute-rs +# hadolint ignore=DL3022 COPY --from=asap-sketchlib . asap_sketchlib +# hadolint ignore=DL3022 COPY --from=asap-gorilla-rust . ASAPCollector/asap-gorilla-rust -RUN cd ASAPQuery-backend && cargo build --release --bin control_plane +# See data_plane/Dockerfile: resolve private sibling dependencies from the +# build contexts instead of requiring Docker to authenticate to GitHub. +RUN if [ "$LOCAL_DEP_OVERRIDES" = true ]; then mkdir -p ASAPQuery-backend/.cargo && \ + printf '%s\n' \ + '[patch."https://github.com/ProjectASAP/ASAPCollector"]' \ + 'asap-precompute-rs = { path = "/src/ASAPCollector/asap-precompute-rs" }' \ + '' \ + '[patch."https://github.com/ProjectASAP/asap_sketchlib"]' \ + 'asap_sketchlib = { path = "/src/asap_sketchlib" }' \ + >> ASAPQuery-backend/.cargo/config.toml; fi + +WORKDIR /src/ASAPQuery-backend +RUN cargo build --release --bin control_plane --bin control_plane_quote_snapshot # Runtime image: binary + CA certs. FROM debian:bookworm-slim @@ -40,6 +56,8 @@ RUN apt-get update && \ rm -rf /var/lib/apt/lists/* COPY --from=build /src/ASAPQuery-backend/target/release/control_plane \ /usr/local/bin/control_plane +COPY --from=build /src/ASAPQuery-backend/target/release/control_plane_quote_snapshot \ + /usr/local/bin/control_plane_quote_snapshot ENV RUST_LOG=info # OpAMP ws 4320, controller gRPC 4321, controller HTTP 8080. diff --git a/control_plane/src/bin/control_plane_quote_snapshot.rs b/control_plane/src/bin/control_plane_quote_snapshot.rs new file mode 100644 index 00000000..e71f21ea --- /dev/null +++ b/control_plane/src/bin/control_plane_quote_snapshot.rs @@ -0,0 +1,95 @@ +//! Complete a backend-local planning snapshot with deterministic unit-cost +//! quotes for the PromQL compliance harness. Candidate identities and demand +//! still come from the production planner and physical compiler. + +use control_plane::physical::{ + compiler::{ + BackendLocalPlanningInput, PhysicalPlanCompiler, BACKEND_REVISION, PLANNER_REVISION, + }, + workload_cost::{ + enumerate_exact_and_materialized_candidates, manifest, WorkloadCostEvidence, WorkloadQuote, + }, +}; +use std::{collections::BTreeMap, path::Path}; + +fn quote_snapshot( + mut snapshot: BackendLocalPlanningInput, +) -> Result { + if snapshot.workload_cost_evidence.is_some() { + return Err("planning snapshot already contains workload_cost_evidence".into()); + } + let observed_at_unix_ms = snapshot.environment.observed_at_unix_ms; + let valid_for_ms = snapshot.environment.max_evidence_age_ms; + let (request, environment) = snapshot + .clone() + .into_physical_compilation_request() + .map_err(|error| error.to_string())?; + let quotes = enumerate_exact_and_materialized_candidates(request) + .map_err(|error| error.to_string())? + .into_iter() + .filter_map(|candidate| { + let plan = PhysicalPlanCompiler + .compile_promql(candidate.clone(), environment.clone()) + .ok()?; + let manifest = manifest(&plan, &candidate.queries).ok()?; + Some(WorkloadQuote { + unit_costs: manifest + .components + .keys() + .map(|key| (key.clone(), 1.0)) + .collect::>(), + manifest, + executable: true, + }) + }) + .collect(); + snapshot.workload_cost_evidence = Some(WorkloadCostEvidence { + backend_revision: BACKEND_REVISION.into(), + planner_revision: PLANNER_REVISION.into(), + data_snapshot_id: "promql-compliance".into(), + model_version: "promql-compliance-unit-costs".into(), + observed_at_unix_ms, + valid_for_ms, + quotes, + }); + Ok(snapshot) +} + +fn main() -> Result<(), String> { + let mut args = std::env::args_os().skip(1); + let input = args + .next() + .ok_or("usage: control_plane_quote_snapshot INPUT OUTPUT")?; + let output = args + .next() + .ok_or("usage: control_plane_quote_snapshot INPUT OUTPUT")?; + if args.next().is_some() { + return Err("usage: control_plane_quote_snapshot INPUT OUTPUT".into()); + } + let snapshot = std::fs::read(&input) + .map_err(|error| format!("read {}: {error}", Path::new(&input).display()))?; + let snapshot = + serde_json::from_slice(&snapshot).map_err(|error| format!("decode snapshot: {error}"))?; + let quoted = quote_snapshot(snapshot)?; + std::fs::write( + &output, + serde_json::to_vec_pretty("ed).map_err(|error| error.to_string())?, + ) + .map_err(|error| format!("write {}: {error}", Path::new(&output).display()))?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn adds_complete_candidate_quotes_to_a_snapshot() { + let fixture = + include_str!("../../../docs/examples/asapquery-compatibility-demo-snapshot.json"); + let mut snapshot: BackendLocalPlanningInput = serde_json::from_str(fixture).unwrap(); + snapshot.workload_cost_evidence = None; + let quoted = quote_snapshot(snapshot).unwrap(); + assert!(!quoted.workload_cost_evidence.unwrap().quotes.is_empty()); + } +} diff --git a/data_plane/Dockerfile b/data_plane/Dockerfile index 7c5ea7b3..cc5e2a72 100644 --- a/data_plane/Dockerfile +++ b/data_plane/Dockerfile @@ -19,6 +19,7 @@ # --build-context asap-gorilla-rust=/path/ASAPCollector/asap-gorilla-rust \ # -t asap/data-plane:dev . FROM rust:1.90-bookworm AS build +ARG LOCAL_DEP_OVERRIDES=false WORKDIR /src # protoc for the backend's prost-build usages; git for any git-dep fetches. @@ -31,11 +32,27 @@ RUN apt-get update && \ # /src/ASAPQuery-backend, /src/ASAPCollector/asap-precompute-rs, # /src/asap_sketchlib, /src/ASAPCollector/asap-gorilla-rust COPY . ASAPQuery-backend +# hadolint ignore=DL3022 COPY --from=asap-precompute-rs . ASAPCollector/asap-precompute-rs +# hadolint ignore=DL3022 COPY --from=asap-sketchlib . asap_sketchlib +# hadolint ignore=DL3022 COPY --from=asap-gorilla-rust . ASAPCollector/asap-gorilla-rust -RUN cd ASAPQuery-backend && cargo build --release --bin data_plane +# Build from the supplied sibling checkouts. CI and local Compose builds do +# not need credentials for ASAPCollector or Sketchlib while this dependency +# remains private. +RUN if [ "$LOCAL_DEP_OVERRIDES" = true ]; then mkdir -p ASAPQuery-backend/.cargo && \ + printf '%s\n' \ + '[patch."https://github.com/ProjectASAP/ASAPCollector"]' \ + 'asap-precompute-rs = { path = "/src/ASAPCollector/asap-precompute-rs" }' \ + '' \ + '[patch."https://github.com/ProjectASAP/asap_sketchlib"]' \ + 'asap_sketchlib = { path = "/src/asap_sketchlib" }' \ + >> ASAPQuery-backend/.cargo/config.toml; fi + +WORKDIR /src/ASAPQuery-backend +RUN cargo build --release --bin data_plane # Runtime image: binary + CA certs. FROM debian:bookworm-slim diff --git a/data_plane/src/main.rs b/data_plane/src/main.rs index 8b106212..ca32bb8a 100644 --- a/data_plane/src/main.rs +++ b/data_plane/src/main.rs @@ -430,11 +430,6 @@ fn validate_profile(args: &Args) -> Result<()> { ) .into()); } - if !args.forward_unsupported_queries { - return Err( - "--profile asapquery requires --forward-unsupported-queries for exact fallback".into(), - ); - } let required_horizon = (args.precompute_allowed_lateness_ms.max(0) as u64) .saturating_add(args.remote_write_expected_retry_interval_ms); if args.remote_write_dedup_horizon_ms < required_horizon { @@ -1503,6 +1498,19 @@ mod tests { .unwrap(); assert!(validate_profile(&valid).is_ok()); + let no_fallback = Args::try_parse_from([ + "data_plane", + "--profile", + "asapquery", + "--physical-plan", + "plan.json", + ]) + .unwrap(); + assert!( + validate_profile(&no_fallback).is_ok(), + "a backend-local plan must be allowed to reject unsupported queries" + ); + let legacy = Args::try_parse_from([ "data_plane", "--profile", diff --git a/docs/design_docs/promql-compliance-grill.md b/docs/design_docs/promql-compliance-grill.md new file mode 100644 index 00000000..61f13173 --- /dev/null +++ b/docs/design_docs/promql-compliance-grill.md @@ -0,0 +1,79 @@ +# PromQL compliance design Q&A + +> **Question:** should v1 be a strict differential suite for a deliberately supported, locally answered subset of PromQL—with unsupported queries required to fail—or should it allow Prometheus fallback? +> +> **Recommendation:** require local answers and reject fallback in v1. Otherwise Prometheus can answer on the backend’s behalf, yielding a green comparison that proves neither backend ingestion nor backend query semantics. + +**Answer:** unsupported queries required to fail, yes + +> **Question:** should the suite compare approximate sketch-backed values to Prometheus within an explicit per-query tolerance, or start with only result shapes/data sizes where equality is exact? +> +> **Recommendation:** support explicit per-query tolerances from day one, defaulting to exact equality. The backend is designed to return approximate sketch results, so pretending all valid local answers are exact would either constrain the suite to an unrepresentative corpus or create noisy failures. + +**Answer:** yes + +> **Question:** should v1 use a finite, timestamp-pinned fixture and explicitly call the backend’s `/api/v1/precompute/drain` before querying? +> +> **Recommendation:** yes. It makes ingestion completion deterministic and avoids flaky “did the asynchronous precompute path catch up?” failures. Live-scrape/streaming behavior should be a separate suite later, because `drain` deliberately closes the backend’s Remote Write input. + +**Answer:** yes + +> **Question:** before the backend can accept Remote Write, it needs a configuration called a PhysicalPlan that says what metrics/queries it should accelerate. Should this new test suite: (1) give the backend a small, fixed test configuration; or (2) start the planner/control plane too and have it create that configuration? +> +> **Recommendation:** option 1 for this issue. It keeps failures clear: if the test fails, it is likely ingestion or query behavior in the backend, not plan generation. + +**Answer:** i wanna do option 2 actually. we can add option 1 later + +> **Question:** should “planner/control plane” mean the backend repository’s `control_plane` binary and its normal plan-publication protocol, rather than importing ASAPQuery’s older `asap-planner-rs` service? +> +> **Recommendation:** use the backend’s own `control_plane`. It tests the deployment users of `ASAPQuery-backend` actually run and avoids coupling this repository’s suite to a separate repository’s planner implementation. + +**Answer:** yes + +> **Question:** should we port the existing Go-based seeder/comparator and Docker Compose runner essentially intact, adapting only the backend service wiring, or rewrite it as Rust process tests? +> +> **Recommendation:** port the Go tooling and Compose workflow. It already provides deterministic Remote Write encoding, structured result comparison, fixtures, tolerances, reports, and lifecycle handling; a rewrite would add risk without improving the test’s signal. + +**Answer:** yes + +> **Question:** what should be the first required green corpus? +> +> **Recommendation:** begin with one small fixture and suite covering a selector, `rate`, `increase`, one aggregation, and both instant/range evaluations—plus one intentionally unsupported query that must fail locally. Add sparse/staleness, counter resets, and larger aggregation matrices after this vertical slice is stable. + +**Answer:** Let's just add all corpora. I think we may not get any green righ tnow + +> **Question:** when those corpora find differences, should the runner exit non-zero from day one? +> +> **Recommendation:** yes—never encode known failures as passes. Keep the full workflow manual/non-required initially, preserve its JSON reports as evidence, and only make a stable subset required in PR CI after the gaps are fixed. + +**Answer:** yes + +> **Question:** should every successful backend response be required to carry and pass a local-execution provenance check (for example, its existing ASAPQuery data-source marker), in addition to disabling fallback? +> +> **Recommendation:** yes. Disabling fallback catches most masking, but provenance makes the test’s claim explicit and will catch accidental routing changes that still return a successful response. + +**Answer:** no, it's fine. just disable fallback + +> **Question:** should the runner wait until the control plane has published and the backend reports an active Remote Write-ready plan before it sends any fixture data? +> +> **Recommendation:** yes. Then seed both targets, call the backend’s finite-input `drain`, and only then execute the fixed-time queries. This makes plan activation and ingestion completion explicit rather than timing-dependent. + +**Answer:** yes + +> **Question:** should the runner derive the control-plane workload/configuration from the same query-suite YAML it executes, rather than maintain a second hand-written plan configuration per corpus? +> +> **Recommendation:** yes. One source of truth prevents a test from querying expressions that the control plane was never asked to plan, and it makes adding a corpus a fixture-only change. + +**Answer:** yes + +> **Question:** on failure, should the runner retain the JSON report and collect service logs, while still tearing down containers by default? +> +> **Recommendation:** yes. CI should upload the report and logs as artifacts; local runs should offer `--keep-services` for interactive debugging. Default cleanup prevents stale volumes/ports from contaminating the next run. + +**Answer:** yes + +> **Question:** should `ASAPQuery-backend` own a copied/adapted version of the harness and fixtures, rather than invoke `ASAPQuery/promql-compliance` across repositories? +> +> **Recommendation:** own it in the backend repository. The backend needs different service wiring—its `control_plane`, PhysicalPlan lifecycle, and one HTTP listener—and an external cross-repo dependency would make local and CI runs less reproducible. + +**Answer:** yes diff --git a/promql-compliance/README.md b/promql-compliance/README.md new file mode 100644 index 00000000..a482a31c --- /dev/null +++ b/promql-compliance/README.md @@ -0,0 +1,98 @@ +# PromQL compliance suite + +This developer harness sends one deterministic Remote Write fixture to Prometheus and ASAPQuery-backend, then compares their Prometheus API responses. It also derives the backend planning snapshot from the same suite queries. + +Run from `promql-compliance/runner`: + +```bash +make run-all +``` + +The first Docker build can take several minutes. Services and volumes are removed after every case. Override artifact locations when needed: + +```bash +make run-all REPORT_DIR=/tmp/promql-reports LOGS_DIR=/tmp/promql-logs +``` + +## Read results + +`REPORT_DIR` defaults to `/tmp/asapquery-backend-promql-reports` and contains: + +- `summary.md`: report card for every dataset/suite case. +- `summary.json`: the same summary for tooling. +- one JSON file per case: every query comparison, raw Prometheus response, raw backend response, and range/instant parity evidence. + +`LOGS_DIR` defaults to `/tmp/asapquery-backend-promql-logs` and contains `compose.log` per case. + +`Overall: true` means every response matched and the backend response was served by ASAPQuery. A response served by `prometheus_fallback` does not pass: matching a forwarded response is not a differential backend result. + +Inspect a concrete query: + +```bash +jq '.queries[] | select(.name == "quantile-over-time-0.5") | .instant[0].responses' \ + /tmp/asapquery-backend-promql-reports/aggregations.json +``` + +The backend payload has `servedBy`: + +- `asap_query` (or another ASAP backend source): answered locally. +- `prometheus_fallback`: forwarded to Prometheus and therefore fails the strict compliance check. + +## Add a query + +Edit or add a YAML suite in `suites/`. Every query needs a stable name, a PromQL expression, and at least one instant offset or a range. Offsets are seconds after the fixture base time. + +```yaml +name: my-suite +comparison_defaults: + value_tolerance: + relative: 0 + absolute: 0.000001 +queries: + - name: median-by-job + expr: quantile_over_time(0.5, data[5m]) + instant_offsets_seconds: [300, 600] + range: + start_offset_seconds: 300 + end_offset_seconds: 600 + step_seconds: 60 +``` + +Run one suite against a fixture: + +```bash +make run DATASET=../datasets/aggregations.yaml SUITE=../suites/my-suite.yaml +``` + +## Change tolerance + +Set a suite default under `comparison_defaults`, or override one query under `comparison`. Relative and absolute finite-value tolerances combine; labels and timestamps always match exactly. + +```yaml +comparison: + value_tolerance: + relative: 0.01 + absolute: 0.000001 +``` + +## Add a dataset + +Add YAML under `datasets/`. Each series has a metric name, labels, and strictly increasing finite samples. Times are offsets in seconds from the run base time. + +```yaml +name: my-data +series: + - metric: data + labels: {job: frontend, instance: i-1} + samples: + - {offset_seconds: 0, value: 10} + - {offset_seconds: 60, value: 20} +``` + +Use a distinct label set for each series of a metric. Ensure query windows have enough samples at every chosen evaluation point. + +## Planning workload and metric vocabulary + +There is no separate hand-written PhysicalPlan. `BuildPlanningSnapshot` derives a backend-local planning snapshot from every suite query; the control-plane helper enumerates candidates and supplies deterministic unit-cost evidence before the data plane starts. + +The suite query expression is the workload query. Dataset metric names are the ingested metric vocabulary and must cover every named metric referenced by a query. Current generated planning defaults are intentionally centralized in `runner/control_plane.go`: 60-second demand cadence for instant-only queries, range step for range queries, explicit 1% epsilon/delta accuracy, declared ingestion rate 100 samples/second, and 1-second scrape interval. Change those defaults there when the desired workload model changes, then run `go test ./...` and the relevant live case. diff --git a/promql-compliance/datasets/aggregations-dense-cadence.yaml b/promql-compliance/datasets/aggregations-dense-cadence.yaml new file mode 100644 index 00000000..f7f719f9 --- /dev/null +++ b/promql-compliance/datasets/aggregations-dense-cadence.yaml @@ -0,0 +1,181 @@ +name: aggregations-dense-cadence +series: + # Same as aggregations.yaml, but every series samples at 60s (the query + # step) instead of backend/i-2 and worker/i-2 sampling slower than it -- + # that mismatch is #705's sparse-cadence bug. aggregations.yaml is left + # as the red case tracking it; this is the green baseline for the same + # query shapes. + # + # Each series also has one extra sample at offset_seconds 1260, past the + # suite's last evaluated offset (1200) so the trailing window has a real + # later sample to close on, instead of only the wall-clock idle fallback. + # + # Known gap (roborev job 191): every series shares the same cadence, so + # count_over_time(data[5m]) ties across all of them. topk-*-count-over-time + # and topk-*-count-over-time-by-job may pick different tied members on + # Prometheus vs. ASAPQuery -- not a reliable green signal for those two + # query shapes specifically. (topk-*-count-over-time-by-job-instance is + # fine: job+instance already uniquely identifies each series, so those + # groups have one member each and can't tie.) + - metric: data + labels: + job: frontend + instance: i-1 + samples: + - {offset_seconds: 0, value: 100} + - {offset_seconds: 60, value: 160} + - {offset_seconds: 120, value: 220} + - {offset_seconds: 180, value: 280} + - {offset_seconds: 240, value: 340} + - {offset_seconds: 300, value: 400} + - {offset_seconds: 360, value: 460} + - {offset_seconds: 420, value: 520} + - {offset_seconds: 480, value: 580} + - {offset_seconds: 540, value: 640} + - {offset_seconds: 600, value: 700} + - {offset_seconds: 660, value: 760} + - {offset_seconds: 720, value: 820} + - {offset_seconds: 780, value: 880} + - {offset_seconds: 840, value: 940} + - {offset_seconds: 900, value: 1000} + - {offset_seconds: 960, value: 1060} + - {offset_seconds: 1020, value: 1120} + - {offset_seconds: 1080, value: 1180} + - {offset_seconds: 1140, value: 1240} + - {offset_seconds: 1200, value: 1300} + - {offset_seconds: 1260, value: 1360} + - metric: data + labels: + job: frontend + instance: i-2 + samples: + - {offset_seconds: 0, value: 200} + - {offset_seconds: 60, value: 320} + - {offset_seconds: 120, value: 440} + - {offset_seconds: 180, value: 560} + - {offset_seconds: 240, value: 680} + - {offset_seconds: 300, value: 800} + - {offset_seconds: 360, value: 920} + - {offset_seconds: 420, value: 1040} + - {offset_seconds: 480, value: 1160} + - {offset_seconds: 540, value: 1280} + - {offset_seconds: 600, value: 1400} + - {offset_seconds: 660, value: 1520} + - {offset_seconds: 720, value: 1640} + - {offset_seconds: 780, value: 1760} + - {offset_seconds: 840, value: 1880} + - {offset_seconds: 900, value: 2000} + - {offset_seconds: 960, value: 2120} + - {offset_seconds: 1020, value: 2240} + - {offset_seconds: 1080, value: 2360} + - {offset_seconds: 1140, value: 2480} + - {offset_seconds: 1200, value: 2600} + - {offset_seconds: 1260, value: 2720} + - metric: data + labels: + job: backend + instance: i-1 + samples: + - {offset_seconds: 0, value: 300} + - {offset_seconds: 60, value: 480} + - {offset_seconds: 120, value: 660} + - {offset_seconds: 180, value: 840} + - {offset_seconds: 240, value: 1020} + - {offset_seconds: 300, value: 1200} + - {offset_seconds: 360, value: 1380} + - {offset_seconds: 420, value: 1560} + - {offset_seconds: 480, value: 1740} + - {offset_seconds: 540, value: 1920} + - {offset_seconds: 600, value: 2100} + - {offset_seconds: 660, value: 2280} + - {offset_seconds: 720, value: 2460} + - {offset_seconds: 780, value: 2640} + - {offset_seconds: 840, value: 2820} + - {offset_seconds: 900, value: 3000} + - {offset_seconds: 960, value: 3180} + - {offset_seconds: 1020, value: 3360} + - {offset_seconds: 1080, value: 3540} + - {offset_seconds: 1140, value: 3720} + - {offset_seconds: 1200, value: 3900} + - {offset_seconds: 1260, value: 4080} + - metric: data + labels: + job: backend + instance: i-2 + samples: + - {offset_seconds: 0, value: 400} + - {offset_seconds: 60, value: 640} + - {offset_seconds: 120, value: 880} + - {offset_seconds: 180, value: 1120} + - {offset_seconds: 240, value: 1360} + - {offset_seconds: 300, value: 1600} + - {offset_seconds: 360, value: 1840} + - {offset_seconds: 420, value: 2080} + - {offset_seconds: 480, value: 2320} + - {offset_seconds: 540, value: 2560} + - {offset_seconds: 600, value: 2800} + - {offset_seconds: 660, value: 3040} + - {offset_seconds: 720, value: 3280} + - {offset_seconds: 780, value: 3520} + - {offset_seconds: 840, value: 3760} + - {offset_seconds: 900, value: 4000} + - {offset_seconds: 960, value: 4240} + - {offset_seconds: 1020, value: 4480} + - {offset_seconds: 1080, value: 4720} + - {offset_seconds: 1140, value: 4960} + - {offset_seconds: 1200, value: 5200} + - {offset_seconds: 1260, value: 5440} + - metric: data + labels: + job: worker + instance: i-1 + samples: + - {offset_seconds: 0, value: 500} + - {offset_seconds: 60, value: 800} + - {offset_seconds: 120, value: 1100} + - {offset_seconds: 180, value: 1400} + - {offset_seconds: 240, value: 1700} + - {offset_seconds: 300, value: 2000} + - {offset_seconds: 360, value: 2300} + - {offset_seconds: 420, value: 2600} + - {offset_seconds: 480, value: 2900} + - {offset_seconds: 540, value: 3200} + - {offset_seconds: 600, value: 3500} + - {offset_seconds: 660, value: 3800} + - {offset_seconds: 720, value: 4100} + - {offset_seconds: 780, value: 4400} + - {offset_seconds: 840, value: 4700} + - {offset_seconds: 900, value: 5000} + - {offset_seconds: 960, value: 5300} + - {offset_seconds: 1020, value: 5600} + - {offset_seconds: 1080, value: 5900} + - {offset_seconds: 1140, value: 6200} + - {offset_seconds: 1200, value: 6500} + - {offset_seconds: 1260, value: 6800} + - metric: data + labels: + job: worker + instance: i-2 + samples: + - {offset_seconds: 0, value: 600} + - {offset_seconds: 60, value: 960} + - {offset_seconds: 120, value: 1320} + - {offset_seconds: 180, value: 1680} + - {offset_seconds: 240, value: 2040} + - {offset_seconds: 300, value: 2400} + - {offset_seconds: 360, value: 2760} + - {offset_seconds: 420, value: 3120} + - {offset_seconds: 480, value: 3480} + - {offset_seconds: 540, value: 3840} + - {offset_seconds: 600, value: 4200} + - {offset_seconds: 660, value: 4560} + - {offset_seconds: 720, value: 4920} + - {offset_seconds: 780, value: 5280} + - {offset_seconds: 840, value: 5640} + - {offset_seconds: 900, value: 6000} + - {offset_seconds: 960, value: 6360} + - {offset_seconds: 1020, value: 6720} + - {offset_seconds: 1080, value: 7080} + - {offset_seconds: 1140, value: 7440} + - {offset_seconds: 1200, value: 7800} + - {offset_seconds: 1260, value: 8160} diff --git a/promql-compliance/datasets/aggregations.yaml b/promql-compliance/datasets/aggregations.yaml new file mode 100644 index 00000000..31c98797 --- /dev/null +++ b/promql-compliance/datasets/aggregations.yaml @@ -0,0 +1,136 @@ +name: aggregations +series: + # The first series in each job is sampled every minute. The second series + # has a sparser cadence so count_over_time and topk have non-trivial output. + - metric: data + labels: + job: frontend + instance: i-1 + samples: + - {offset_seconds: 0, value: 100} + - {offset_seconds: 60, value: 160} + - {offset_seconds: 120, value: 220} + - {offset_seconds: 180, value: 280} + - {offset_seconds: 240, value: 340} + - {offset_seconds: 300, value: 400} + - {offset_seconds: 360, value: 460} + - {offset_seconds: 420, value: 520} + - {offset_seconds: 480, value: 580} + - {offset_seconds: 540, value: 640} + - {offset_seconds: 600, value: 700} + - {offset_seconds: 660, value: 760} + - {offset_seconds: 720, value: 820} + - {offset_seconds: 780, value: 880} + - {offset_seconds: 840, value: 940} + - {offset_seconds: 900, value: 1000} + - {offset_seconds: 960, value: 1060} + - {offset_seconds: 1020, value: 1120} + - {offset_seconds: 1080, value: 1180} + - {offset_seconds: 1140, value: 1240} + - {offset_seconds: 1200, value: 1300} + - metric: data + labels: + job: frontend + instance: i-2 + samples: + - {offset_seconds: 0, value: 200} + - {offset_seconds: 60, value: 320} + - {offset_seconds: 120, value: 440} + - {offset_seconds: 180, value: 560} + - {offset_seconds: 240, value: 680} + - {offset_seconds: 300, value: 800} + - {offset_seconds: 360, value: 920} + - {offset_seconds: 420, value: 1040} + - {offset_seconds: 480, value: 1160} + - {offset_seconds: 540, value: 1280} + - {offset_seconds: 600, value: 1400} + - {offset_seconds: 660, value: 1520} + - {offset_seconds: 720, value: 1640} + - {offset_seconds: 780, value: 1760} + - {offset_seconds: 840, value: 1880} + - {offset_seconds: 900, value: 2000} + - {offset_seconds: 960, value: 2120} + - {offset_seconds: 1020, value: 2240} + - {offset_seconds: 1080, value: 2360} + - {offset_seconds: 1140, value: 2480} + - {offset_seconds: 1200, value: 2600} + - metric: data + labels: + job: backend + instance: i-1 + samples: + - {offset_seconds: 0, value: 300} + - {offset_seconds: 60, value: 480} + - {offset_seconds: 120, value: 660} + - {offset_seconds: 180, value: 840} + - {offset_seconds: 240, value: 1020} + - {offset_seconds: 300, value: 1200} + - {offset_seconds: 360, value: 1380} + - {offset_seconds: 420, value: 1560} + - {offset_seconds: 480, value: 1740} + - {offset_seconds: 540, value: 1920} + - {offset_seconds: 600, value: 2100} + - {offset_seconds: 660, value: 2280} + - {offset_seconds: 720, value: 2460} + - {offset_seconds: 780, value: 2640} + - {offset_seconds: 840, value: 2820} + - {offset_seconds: 900, value: 3000} + - {offset_seconds: 960, value: 3180} + - {offset_seconds: 1020, value: 3360} + - {offset_seconds: 1080, value: 3540} + - {offset_seconds: 1140, value: 3720} + - {offset_seconds: 1200, value: 3900} + - metric: data + labels: + job: backend + instance: i-2 + samples: + - {offset_seconds: 0, value: 400} + - {offset_seconds: 120, value: 880} + - {offset_seconds: 240, value: 1360} + - {offset_seconds: 360, value: 1840} + - {offset_seconds: 480, value: 2320} + - {offset_seconds: 600, value: 2800} + - {offset_seconds: 720, value: 3280} + - {offset_seconds: 840, value: 3760} + - {offset_seconds: 960, value: 4240} + - {offset_seconds: 1080, value: 4720} + - {offset_seconds: 1200, value: 5200} + - metric: data + labels: + job: worker + instance: i-1 + samples: + - {offset_seconds: 0, value: 500} + - {offset_seconds: 60, value: 800} + - {offset_seconds: 120, value: 1100} + - {offset_seconds: 180, value: 1400} + - {offset_seconds: 240, value: 1700} + - {offset_seconds: 300, value: 2000} + - {offset_seconds: 360, value: 2300} + - {offset_seconds: 420, value: 2600} + - {offset_seconds: 480, value: 2900} + - {offset_seconds: 540, value: 3200} + - {offset_seconds: 600, value: 3500} + - {offset_seconds: 660, value: 3800} + - {offset_seconds: 720, value: 4100} + - {offset_seconds: 780, value: 4400} + - {offset_seconds: 840, value: 4700} + - {offset_seconds: 900, value: 5000} + - {offset_seconds: 960, value: 5300} + - {offset_seconds: 1020, value: 5600} + - {offset_seconds: 1080, value: 5900} + - {offset_seconds: 1140, value: 6200} + - {offset_seconds: 1200, value: 6500} + - metric: data + labels: + job: worker + instance: i-2 + samples: + - {offset_seconds: 0, value: 600} + - {offset_seconds: 180, value: 1680} + - {offset_seconds: 360, value: 2760} + - {offset_seconds: 540, value: 3840} + - {offset_seconds: 720, value: 4920} + - {offset_seconds: 900, value: 6000} + - {offset_seconds: 1080, value: 7080} diff --git a/promql-compliance/datasets/single-rate.yaml b/promql-compliance/datasets/single-rate.yaml new file mode 100644 index 00000000..8bc52bf9 --- /dev/null +++ b/promql-compliance/datasets/single-rate.yaml @@ -0,0 +1,27 @@ +name: single-rate +series: + - metric: http_requests_total + labels: + host: a + samples: + - {offset_seconds: 0, value: 0} + - {offset_seconds: 60, value: 60} + - {offset_seconds: 120, value: 120} + - {offset_seconds: 180, value: 180} + - {offset_seconds: 240, value: 240} + - {offset_seconds: 300, value: 300} + - {offset_seconds: 360, value: 360} + - {offset_seconds: 420, value: 420} + - {offset_seconds: 480, value: 480} + - {offset_seconds: 540, value: 540} + - {offset_seconds: 600, value: 600} + - {offset_seconds: 660, value: 660} + - {offset_seconds: 720, value: 720} + - {offset_seconds: 780, value: 780} + - {offset_seconds: 840, value: 840} + - {offset_seconds: 900, value: 900} + - {offset_seconds: 960, value: 960} + - {offset_seconds: 1020, value: 1020} + - {offset_seconds: 1080, value: 1080} + - {offset_seconds: 1140, value: 1140} + - {offset_seconds: 1200, value: 1200} diff --git a/promql-compliance/datasets/sparse-checkout.yaml b/promql-compliance/datasets/sparse-checkout.yaml new file mode 100644 index 00000000..d9b71da8 --- /dev/null +++ b/promql-compliance/datasets/sparse-checkout.yaml @@ -0,0 +1,74 @@ +name: sparse-checkout +series: + - metric: http_requests_total + labels: + host: a + samples: + - {offset_seconds: 0, value: 0} + - {offset_seconds: 60, value: 60} + - {offset_seconds: 120, value: 120} + - {offset_seconds: 180, value: 180} + - {offset_seconds: 240, value: 240} + - {offset_seconds: 300, value: 300} + - {offset_seconds: 360, value: 360} + - {offset_seconds: 420, value: 420} + - {offset_seconds: 480, value: 480} + - {offset_seconds: 540, value: 540} + - {offset_seconds: 600, value: 600} + - {offset_seconds: 660, value: 660} + - {offset_seconds: 720, value: 720} + - {offset_seconds: 780, value: 780} + - {offset_seconds: 840, value: 840} + - {offset_seconds: 900, value: 900} + - {offset_seconds: 960, value: 960} + - {offset_seconds: 1020, value: 1020} + - {offset_seconds: 1080, value: 1080} + - {offset_seconds: 1140, value: 1140} + - {offset_seconds: 1200, value: 1200} + - metric: http_requests_total + labels: + host: b + samples: + - {offset_seconds: 0, value: 1000} + - {offset_seconds: 60, value: 1120} + - {offset_seconds: 120, value: 1240} + - {offset_seconds: 180, value: 1360} + - {offset_seconds: 240, value: 1480} + - {offset_seconds: 300, value: 1600} + - {offset_seconds: 360, value: 1720} + - {offset_seconds: 420, value: 1840} + - {offset_seconds: 480, value: 1960} + - {offset_seconds: 540, value: 2080} + - {offset_seconds: 600, value: 2200} + - {offset_seconds: 660, value: 2320} + - {offset_seconds: 720, value: 2440} + - {offset_seconds: 780, value: 2560} + - {offset_seconds: 840, value: 2680} + - {offset_seconds: 900, value: 2800} + - {offset_seconds: 960, value: 2920} + - {offset_seconds: 1020, value: 3040} + - {offset_seconds: 1080, value: 3160} + - {offset_seconds: 1140, value: 3280} + - {offset_seconds: 1200, value: 3400} + - metric: checkout_up + labels: + service: checkout + region: us-east + samples: + - {offset_seconds: 0, value: 1} + - {offset_seconds: 60, value: 1} + - {offset_seconds: 120, value: 1} + - {offset_seconds: 180, value: 1} + - {offset_seconds: 240, value: 1} + - {offset_seconds: 300, value: 1} + - metric: checkout_up + labels: + service: checkout + region: us-west + samples: + - {offset_seconds: 900, value: 1} + - {offset_seconds: 960, value: 1} + - {offset_seconds: 1020, value: 1} + - {offset_seconds: 1080, value: 1} + - {offset_seconds: 1140, value: 1} + - {offset_seconds: 1200, value: 1} diff --git a/promql-compliance/docker-compose.yml b/promql-compliance/docker-compose.yml new file mode 100644 index 00000000..64f1b082 --- /dev/null +++ b/promql-compliance/docker-compose.yml @@ -0,0 +1,34 @@ +name: asapquery-backend-promql-compliance + +services: + prometheus: + image: prom/prometheus:v3.9.1 + ports: ["${PROMETHEUS_PORT:-19090}:9090"] + command: ["--config.file=/etc/prometheus/prometheus.yml", "--web.enable-remote-write-receiver", "--storage.tsdb.path=/prometheus"] + planner: + build: + context: .. + dockerfile: control_plane/Dockerfile + args: {LOCAL_DEP_OVERRIDES: "true"} + additional_contexts: + asap-precompute-rs: ${ASAP_PRECOMPUTE_RS_CONTEXT} + asap-sketchlib: ${ASAP_SKETCHLIB_CONTEXT} + asap-gorilla-rust: ${ASAP_GORILLA_RUST_CONTEXT} + entrypoint: ["/usr/local/bin/control_plane_quote_snapshot"] + command: ["/config/planning-snapshot-template.json", "/config/planning-snapshot.json"] + volumes: + - ${ASAP_PLANNING_SNAPSHOT_TEMPLATE}:/config/planning-snapshot-template.json:ro + - ${ASAP_PLANNING_SNAPSHOT}:/config/planning-snapshot.json + data-plane: + build: + context: .. + dockerfile: data_plane/Dockerfile + args: {LOCAL_DEP_OVERRIDES: "true"} + additional_contexts: + asap-precompute-rs: ${ASAP_PRECOMPUTE_RS_CONTEXT} + asap-sketchlib: ${ASAP_SKETCHLIB_CONTEXT} + asap-gorilla-rust: ${ASAP_GORILLA_RUST_CONTEXT} + command: ["--profile", "asapquery", "--planning-snapshot", "/config/planning-snapshot.json", "--prometheus-server", "http://prometheus:9090", "--forward-unsupported-queries", "--http-port", "9091", "--output-dir", "/tmp/asap"] + volumes: + - ${ASAP_PLANNING_SNAPSHOT}:/config/planning-snapshot.json:ro + ports: ["${BACKEND_PORT:-19091}:9091"] diff --git a/promql-compliance/runner/Makefile b/promql-compliance/runner/Makefile new file mode 100644 index 00000000..1c4b4985 --- /dev/null +++ b/promql-compliance/runner/Makefile @@ -0,0 +1,29 @@ +.PHONY: test run run-all + +DATASET ?= ../datasets/single-rate.yaml +SUITE ?= ../suites/temporal.yaml +REPORT_DIR ?= /tmp/asapquery-backend-promql-reports +REFERENCE_URL ?= http://localhost:19090 +BACKEND_URL ?= http://localhost:19091 +COMPOSE_FILE ?= ../docker-compose.yml +LOGS_DIR ?= /tmp/asapquery-backend-promql-logs + +CASES := \ + single-rate-temporal:../datasets/single-rate.yaml:../suites/temporal.yaml \ + sparse-checkout-temporal:../datasets/sparse-checkout.yaml:../suites/temporal.yaml \ + aggregations:../datasets/aggregations.yaml:../suites/aggregations.yaml \ + aggregations-dense-cadence:../datasets/aggregations-dense-cadence.yaml:../suites/aggregations.yaml + +test: + go test ./... + +run: + go run ./cmd/differential-runner --dataset $(DATASET) --suite $(SUITE) --reference-url $(REFERENCE_URL) --test-url $(BACKEND_URL) --compose-file $(COMPOSE_FILE) --logs-dir $(LOGS_DIR) + +run-all: + @set -eu; mkdir -p "$(REPORT_DIR)"; rm -f "$(REPORT_DIR)"/*.json "$(REPORT_DIR)"/summary.md; result=0; \ + for test_case in $(CASES); do \ + name="$${test_case%%:*}"; rest="$${test_case#*:}"; dataset="$${rest%%:*}"; suite="$${rest#*:}"; \ + echo "Running $$name"; \ + if ! go run ./cmd/differential-runner --dataset "$$dataset" --suite "$$suite" --reference-url "$(REFERENCE_URL)" --test-url "$(BACKEND_URL)" --compose-file "$(COMPOSE_FILE)" --logs-dir "$(LOGS_DIR)/$$name" --output "$(REPORT_DIR)/$$name.json"; then result=1; fi; \ + done; go run ./cmd/report-card --reports-dir "$(REPORT_DIR)"; exit $$result diff --git a/promql-compliance/runner/cmd/differential-runner/main.go b/promql-compliance/runner/cmd/differential-runner/main.go new file mode 100644 index 00000000..e566ab71 --- /dev/null +++ b/promql-compliance/runner/cmd/differential-runner/main.go @@ -0,0 +1,129 @@ +package main + +import ( + "context" + "encoding/json" + "flag" + "fmt" + "log" + "os" + "path/filepath" + "time" + + "github.com/ProjectASAP/ASAPQuery-backend/promql-compliance/runner" +) + +type stringList []string + +func (items *stringList) String() string { return fmt.Sprint([]string(*items)) } +func (items *stringList) Set(value string) error { *items = append(*items, value); return nil } + +func main() { + if err := run(); err != nil { + fmt.Fprintln(os.Stderr, err) + os.Exit(1) + } +} + +func run() error { + var composeFiles stringList + datasetPath := flag.String("dataset", "", "dataset YAML") + suitePath := flag.String("suite", "", "query suite YAML") + reference := flag.String("reference-url", "", "Prometheus URL") + test := flag.String("test-url", "", "backend URL") + baseMillis := flag.Int64("base-time-ms", time.Now().Add(-30*time.Minute).UnixMilli(), "fixture base time") + reportPath := flag.String("output", "differential-report.json", "JSON report path") + composeProject := flag.String("compose-project", "asapquery-backend-promql-compliance", "Compose project name") + logsDirectory := flag.String("logs-dir", "", "directory for retained Compose logs") + keepServices := flag.Bool("keep-services", false, "leave Compose services running after the run") + flag.Var(&composeFiles, "compose-file", "Compose file to start; may be repeated") + flag.Parse() + if *datasetPath == "" || *suitePath == "" || *reference == "" || *test == "" { + flag.Usage() + return fmt.Errorf("dataset, suite, reference-url, and test-url are required") + } + dataset, err := runner.LoadDatasetFile(*datasetPath) + if err != nil { + return err + } + suite, err := runner.LoadSuiteFile(*suitePath) + if err != nil { + return err + } + body, err := runner.EncodeRemoteWrite(*baseMillis, dataset) + if err != nil { + return err + } + ctx := context.Background() + log.Printf("loaded dataset %q and suite %q", dataset.Name, suite.Name) + runDirectory, err := os.MkdirTemp("", "asapquery-promql-compliance-") + if err != nil { + return err + } + defer os.RemoveAll(runDirectory) + snapshot, err := json.Marshal(runner.BuildPlanningSnapshot(suite, time.Now().UTC())) + if err != nil { + return err + } + snapshotTemplatePath := filepath.Join(runDirectory, "planning-snapshot-template.json") + if err := os.WriteFile(snapshotTemplatePath, snapshot, 0o600); err != nil { + return err + } + snapshotPath := filepath.Join(runDirectory, "planning-snapshot.json") + // Docker creates a directory when a bind-mounted file target is absent. + // The planner helper overwrites this placeholder with complete cost evidence. + if err := os.WriteFile(snapshotPath, nil, 0o600); err != nil { + return err + } + log.Printf("wrote suite-derived planning snapshot template to %s", snapshotTemplatePath) + lifecycle := runner.ComposeLifecycle{Files: composeFiles, Project: *composeProject, LogsDirectory: *logsDirectory, PlanningSnapshot: snapshotPath, PlanningSnapshotTemplate: snapshotTemplatePath} + if len(composeFiles) > 0 { + log.Printf("building and starting Compose services; the first run may take several minutes") + } + if len(composeFiles) > 0 && !*keepServices { + defer lifecycle.Stop() + } + if err := lifecycle.Start(ctx); err != nil { + return err + } + log.Printf("waiting for Prometheus at %s", *reference) + if err := runner.WaitForHTTP(ctx, *reference+"/api/v1/status/runtimeinfo"); err != nil { + return err + } + log.Printf("waiting for backend at %s", *test) + if err := runner.WaitForHTTP(ctx, *test+"/api/v1/health"); err != nil { + return err + } + log.Printf("seeding identical Remote Write payloads") + if err := runner.PushRemoteWrite(ctx, body, *reference, *test); err != nil { + return err + } + log.Printf("draining backend precompute work") + if err := runner.Drain(ctx, *test); err != nil { + return err + } + base := time.UnixMilli(*baseMillis) + refTarget, testTarget := runner.HTTPQueryTarget{BaseURL: *reference}, runner.HTTPQueryTarget{BaseURL: *test, BackendTarget: true} + log.Printf("comparing %d query cases", len(suite.Queries)) + report := runner.CompareSuite(ctx, refTarget, testTarget, suite, dataset.Name, base) + if err := os.MkdirAll(filepath.Dir(*reportPath), 0o755); err != nil && filepath.Dir(*reportPath) != "." { + return err + } + file, err := os.Create(*reportPath) + if err != nil { + return err + } + if err := json.NewEncoder(file).Encode(report); err != nil { + _ = file.Close() + return err + } + if err := file.Close(); err != nil { + return err + } + log.Printf("wrote comparison report to %s", *reportPath) + fmt.Printf("dataset=%s suite=%s passed=%t report=%s\n", report.Dataset, report.Suite, report.Passed, *reportPath) + if !report.Passed { + return fmt.Errorf("PromQL comparison failed; report: %s", *reportPath) + } + return nil +} diff --git a/promql-compliance/runner/cmd/report-card/main.go b/promql-compliance/runner/cmd/report-card/main.go new file mode 100644 index 00000000..0db854e6 --- /dev/null +++ b/promql-compliance/runner/cmd/report-card/main.go @@ -0,0 +1,71 @@ +package main + +import ( + "encoding/json" + "flag" + "fmt" + "os" + "path/filepath" + + "github.com/ProjectASAP/ASAPQuery-backend/promql-compliance/runner" +) + +type card struct { + Cases []caseCard `json:"cases"` + Passed bool `json:"passed"` +} +type caseCard struct { + Dataset, Suite string + Passed bool + Queries, ASAPQuery, PrometheusFallback int +} + +func main() { + dir := flag.String("reports-dir", "/tmp/asapquery-backend-promql-reports", "report directory") + flag.Parse() + files, err := filepath.Glob(filepath.Join(*dir, "*.json")) + if err != nil { + panic(err) + } + out := card{Passed: true} + for _, file := range files { + if filepath.Base(file) == "summary.json" { + continue + } + var r runner.Report + raw, e := os.ReadFile(file) + if e != nil { + panic(e) + } + if e = json.Unmarshal(raw, &r); e != nil { + panic(e) + } + c := caseCard{Dataset: r.Dataset, Suite: r.Suite, Passed: r.Passed, Queries: len(r.Queries)} + for _, q := range r.Queries { + if q.RangeResponses != nil { + count(&c, q.RangeResponses.Backend.ServedBy) + } + for _, i := range q.Instant { + count(&c, i.Responses.Backend.ServedBy) + } + } + out.Cases = append(out.Cases, c) + out.Passed = out.Passed && c.Passed + } + b, _ := json.MarshalIndent(out, "", " ") + _ = os.WriteFile(filepath.Join(*dir, "summary.json"), b, 0644) + f, _ := os.Create(filepath.Join(*dir, "summary.md")) + defer f.Close() + fmt.Fprintf(f, "# PromQL compliance report card\n\nOverall: **%t**\n\n| Dataset | Suite | Passed | Queries | ASAPQuery answers | Prometheus fallback |\n|---|---|---:|---:|---:|---:|\n", out.Passed) + for _, c := range out.Cases { + fmt.Fprintf(f, "| %s | %s | %t | %d | %d | %d |\n", c.Dataset, c.Suite, c.Passed, c.Queries, c.ASAPQuery, c.PrometheusFallback) + } +} + +func count(c *caseCard, servedBy string) { + if servedBy == "prometheus_fallback" || servedBy == "" { + c.PrometheusFallback++ + return + } + c.ASAPQuery++ +} diff --git a/promql-compliance/runner/compare.go b/promql-compliance/runner/compare.go new file mode 100644 index 00000000..97d50bac --- /dev/null +++ b/promql-compliance/runner/compare.go @@ -0,0 +1,222 @@ +package runner + +import ( + "encoding/json" + "fmt" + "math" + "sort" + "strconv" + "strings" + "time" +) + +// CompareResponses compares successful Prometheus API responses by result +// semantics. Labels and timestamps are exact; only finite sample values use +// the configured tolerance. +func CompareResponses(reference, test QueryResponse, policy ComparisonPolicy) error { + if reference.Status != "success" || test.Status != "success" { + return fmt.Errorf("query status differs or failed: reference=%q/%q test=%q/%q", reference.Status, reference.Error, test.Status, test.Error) + } + left, err := normalizeResponse(reference) + if err != nil { + return fmt.Errorf("normalize reference response: %w", err) + } + right, err := normalizeResponse(test) + if err != nil { + return fmt.Errorf("normalize test response: %w", err) + } + return compareNormalized(left, right, policy) +} + +type normalizedResult struct { + Type string + Samples []normalizedSample + Scalar *float64 + Text *string +} +type normalizedSample struct { + Labels string + Timestamp int64 + Value float64 +} +type apiResult struct { + ResultType string `json:"resultType"` + Result json.RawMessage `json:"result"` +} +type apiSeries struct { + Metric map[string]string `json:"metric"` + Value json.RawMessage `json:"value"` + Values []json.RawMessage `json:"values"` +} + +func normalizeResponse(response QueryResponse) (normalizedResult, error) { + var data apiResult + if err := json.Unmarshal(response.Data, &data); err != nil { + return normalizedResult{}, err + } + switch data.ResultType { + case "vector", "matrix": + var series []apiSeries + if err := json.Unmarshal(data.Result, &series); err != nil { + return normalizedResult{}, err + } + out := normalizedResult{Type: data.ResultType} + for _, item := range series { + values := item.Values + if data.ResultType == "vector" { + values = []json.RawMessage{item.Value} + } + for _, value := range values { + ts, number, err := parsePoint(value) + if err != nil { + return normalizedResult{}, err + } + out.Samples = append(out.Samples, normalizedSample{Labels: canonicalResponseLabels(item.Metric), Timestamp: ts, Value: number}) + } + } + sort.Slice(out.Samples, func(i, j int) bool { + if out.Samples[i].Labels != out.Samples[j].Labels { + return out.Samples[i].Labels < out.Samples[j].Labels + } + return out.Samples[i].Timestamp < out.Samples[j].Timestamp + }) + return out, nil + case "scalar": + _, value, err := parsePoint(data.Result) + if err != nil { + return normalizedResult{}, err + } + return normalizedResult{Type: data.ResultType, Scalar: &value}, nil + case "string": + var point []json.RawMessage + if err := json.Unmarshal(data.Result, &point); err != nil { + return normalizedResult{}, err + } + if len(point) != 2 { + return normalizedResult{}, fmt.Errorf("string result has %d fields", len(point)) + } + var text string + if err := json.Unmarshal(point[1], &text); err != nil { + return normalizedResult{}, err + } + return normalizedResult{Type: data.ResultType, Text: &text}, nil + default: + return normalizedResult{}, fmt.Errorf("unsupported result type %q", data.ResultType) + } +} + +func parsePoint(raw json.RawMessage) (int64, float64, error) { + var point []json.RawMessage + if err := json.Unmarshal(raw, &point); err != nil { + return 0, 0, err + } + if len(point) != 2 { + return 0, 0, fmt.Errorf("sample has %d fields", len(point)) + } + var timestamp float64 + if err := json.Unmarshal(point[0], ×tamp); err != nil { + return 0, 0, err + } + var text string + if err := json.Unmarshal(point[1], &text); err != nil { + return 0, 0, err + } + value, err := strconv.ParseFloat(text, 64) + if err != nil { + return 0, 0, err + } + return int64(math.Round(timestamp * 1000)), value, nil +} + +func canonicalResponseLabels(labels map[string]string) string { + keys := make([]string, 0, len(labels)) + for key := range labels { + keys = append(keys, key) + } + sort.Strings(keys) + var b strings.Builder + for _, key := range keys { + fmt.Fprintf(&b, "%s=%q,", key, labels[key]) + } + return b.String() +} + +func compareNormalized(left, right normalizedResult, policy ComparisonPolicy) error { + if left.Type != right.Type { + return fmt.Errorf("result type differs: reference=%s test=%s", left.Type, right.Type) + } + if left.Scalar != nil || right.Scalar != nil { + if left.Scalar == nil || right.Scalar == nil || !equalFloat(*left.Scalar, *right.Scalar, policy.ValueTolerance) { + return fmt.Errorf("scalar differs: reference=%v test=%v", left.Scalar, right.Scalar) + } + return nil + } + if left.Text != nil || right.Text != nil { + if left.Text == nil || right.Text == nil || *left.Text != *right.Text { + return fmt.Errorf("string differs: reference=%v test=%v", left.Text, right.Text) + } + return nil + } + if len(left.Samples) != len(right.Samples) { + return fmt.Errorf("sample count differs: reference=%d test=%d", len(left.Samples), len(right.Samples)) + } + for i := range left.Samples { + a, b := left.Samples[i], right.Samples[i] + if a.Labels != b.Labels || a.Timestamp != b.Timestamp { + return fmt.Errorf("sample %d labels or timestamp differs: reference=%+v test=%+v", i, a, b) + } + if !equalFloat(a.Value, b.Value, policy.ValueTolerance) { + return fmt.Errorf("sample %d value differs: reference=%v test=%v", i, a.Value, b.Value) + } + } + return nil +} + +func equalFloat(left, right float64, tolerance *Tolerance) bool { + if math.IsNaN(left) || math.IsNaN(right) { + return math.IsNaN(left) && math.IsNaN(right) + } + if math.IsInf(left, 0) || math.IsInf(right, 0) { + return left == right + } + relative, absolute := 0.0, 0.0 + if tolerance != nil { + if tolerance.Relative != nil { + relative = *tolerance.Relative + } + if tolerance.Absolute != nil { + absolute = *tolerance.Absolute + } + } + return math.Abs(left-right) <= absolute+relative*math.Max(math.Abs(left), math.Abs(right)) +} + +// CompareRangeAtInstant checks that range-at-t and instant-at-t agree. +func CompareRangeAtInstant(rangeResponse, instantResponse QueryResponse, at time.Time, policy ComparisonPolicy) error { + rangeValue, err := normalizeResponse(rangeResponse) + if err != nil { + return err + } + instantValue, err := normalizeResponse(instantResponse) + if err != nil { + return err + } + if rangeValue.Type != "matrix" { + return fmt.Errorf("range response type is %q, want matrix", rangeValue.Type) + } + rangeValue.Type = "vector" + wanted := at.UnixMilli() + filtered := rangeValue.Samples[:0] + for _, sample := range rangeValue.Samples { + if sample.Timestamp == wanted { + filtered = append(filtered, sample) + } + } + rangeValue.Samples = filtered + if instantValue.Type == "vector" { + for i := range instantValue.Samples { + instantValue.Samples[i].Timestamp = wanted + } + } + return compareNormalized(rangeValue, instantValue, policy) +} diff --git a/promql-compliance/runner/compare_test.go b/promql-compliance/runner/compare_test.go new file mode 100644 index 00000000..5d2e2d7d --- /dev/null +++ b/promql-compliance/runner/compare_test.go @@ -0,0 +1,33 @@ +package runner + +import ( + "encoding/json" + "testing" + "time" +) + +func comparisonResponse(t *testing.T, resultType, result string) QueryResponse { + t.Helper() + return QueryResponse{Status: "success", Data: json.RawMessage(`{"resultType":"` + resultType + `","result":` + result + `}`)} +} + +func TestCompareResponsesHonorsValueToleranceButNotLabels(t *testing.T) { + left := comparisonResponse(t, "vector", `[{"metric":{"job":"api"},"value":[1,"100"]}]`) + right := comparisonResponse(t, "vector", `[{"metric":{"job":"api"},"value":[1,"101"]}]`) + relative := 0.02 + if err := CompareResponses(left, right, ComparisonPolicy{ValueTolerance: &Tolerance{Relative: &relative}}); err != nil { + t.Fatalf("CompareResponses: %v", err) + } + differentLabels := comparisonResponse(t, "vector", `[{"metric":{"job":"worker"},"value":[1,"101"]}]`) + if err := CompareResponses(left, differentLabels, ComparisonPolicy{ValueTolerance: &Tolerance{Relative: &relative}}); err == nil { + t.Fatal("comparison accepted different labels") + } +} + +func TestCompareRangeAtInstantChecksRequestedGridPoint(t *testing.T) { + rangeResponse := comparisonResponse(t, "matrix", `[{"metric":{"job":"api"},"values":[[1,"1"],[2,"2"]]}]`) + instantResponse := comparisonResponse(t, "vector", `[{"metric":{"job":"api"},"value":[2,"2"]}]`) + if err := CompareRangeAtInstant(rangeResponse, instantResponse, time.Unix(2, 0), ComparisonPolicy{}); err != nil { + t.Fatalf("CompareRangeAtInstant: %v", err) + } +} diff --git a/promql-compliance/runner/control_plane.go b/promql-compliance/runner/control_plane.go new file mode 100644 index 00000000..f341b70f --- /dev/null +++ b/promql-compliance/runner/control_plane.go @@ -0,0 +1,198 @@ +package runner + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "os" + "os/exec" + "path/filepath" + "regexp" + "strings" + "time" +) + +// BuildPublicationRequest derives the controller workload from the same suite +// that will be evaluated. The finite fixture supplies the metric vocabulary; +// no second hand-written planning configuration is maintained. +func BuildPublicationRequest(dataset Dataset, suite Suite, now time.Time) (map[string]any, error) { + plannerRevision, backendRevision, err := buildRevisions() + if err != nil { + return nil, err + } + metrics := make(map[string]struct{}, len(dataset.Series)) + for _, series := range dataset.Series { + metrics[series.Metric] = struct{}{} + } + queries := make([]any, 0, len(suite.Queries)) + for index, query := range suite.Queries { + metric, ok := metricInQuery(query.Expr, metrics) + if !ok { + return nil, fmt.Errorf("query %q does not reference a fixture metric", query.Name) + } + window := queryWindowSeconds(query.Expr) + interval := uint32(60_000) + if query.Range != nil { + interval = uint32(query.Range.StepSeconds * 1000) + } + queries = append(queries, map[string]any{ + "query_id": fmt.Sprintf("%s-%d", suite.Name, index+1), "query_string": query.Expr, "metric": metric, "window_secs": window, "group_by": []string{}, "accuracy": map[string]any{"Epsilon": 0.01}, "evaluation_phase_ms": 0, + "window_cost_model": map[string]any{"implementation_id": "promql-compliance", "cost": map[string]any{"model_version": "promql-compliance-v1", "workload_fingerprint": query.Name, "observed_at_unix_ms": now.UnixMilli(), "valid_for_ms": 600_000, "horizon_seconds": 300.0, "cpu_cost": 1.0, "weighted_cost": 1.0, "peak_memory_bytes": 4096, "network_bytes": 0, "storage_bytes": 2048, "source_scan_bytes": 0}}, + "lifecycle": map[string]any{"evaluation_interval_ms": interval, "ingestion_rate_per_second": 100.0, "evidence_observed_at_unix_ms": now.UnixMilli(), "evidence_valid_for_ms": 600_000, "horizon_seconds": 300.0, "costs": map[string]any{"build": 10.0, "maintenance_per_update": 0.001, "read": 0.1, "retention_per_second": 0.001, "retirement": 1.0}}, + }) + } + return map[string]any{"target": "backend_local_remote_write", "queries": queries, "collector_ids": []string{}, "capability_snapshot_id": "promql-compliance", "planner_revision": plannerRevision, "max_evidence_age_ms": 600_000, "plan_version": 1, "activation_unix_ms": now.UnixMilli(), "backend_compat": "asap-query-backend.v1", "apply_timeout_ms": 30_000, "backend_revision": backendRevision}, nil +} + +// BuildPlanningSnapshot is the backend-local startup input. The data plane +// invokes the repository's pinned Planner and PhysicalPlanCompiler from this +// workload, so the fixture has one source of truth for planning and queries. +func BuildPlanningSnapshot(suite Suite, now time.Time) map[string]any { + queries := make([]any, 0, len(suite.Queries)) + for _, query := range suite.Queries { + interval := 60_000.0 + if query.Range != nil { + interval = query.Range.StepSeconds * 1000 + } + queries = append(queries, map[string]any{ + "query": query.Expr, + "demand": map[string]any{"fixed_interval_at": map[string]any{"interval": interval, "evaluation_phase": 0}}, + "requirements": map[string]any{"accuracy": map[string]any{"explicit": map[string]any{"EpsilonDelta": map[string]any{"epsilon": 0.01, "delta": 0.01}}}, "response_latency": "unspecified"}, + "predictability": map[string]any{"predictable": map[string]any{"known_at": nil}}, + "time_selection": map[string]any{"scope": "real_time", "lookback": nil, "as_of": nil}, + }) + } + data := map[string]any{"arrival": "continuously_ingesting", "ingestion_volume": map[string]any{"value": nil, "source": "unknown", "observed_at_ms": nil, "valid_for_ms": nil}, "ingestion_rate": map[string]any{"value": 100.0, "source": "declared", "observed_at_ms": nil, "valid_for_ms": nil}, "input_cardinality": map[string]any{"value": nil, "source": "unknown", "observed_at_ms": nil, "valid_for_ms": nil}, "distribution": map[string]any{"value": nil, "source": "unknown", "observed_at_ms": nil, "valid_for_ms": nil}} + return map[string]any{"snapshot_version": 2, "query_workload": map[string]any{"language": "promql", "query_batch": nil, "repeating_queries": queries, "data_workload": data}, "data_workload": data, "implementation": map[string]any{"lifecycle_costs": map[string]any{"build": 10.0, "maintenance_per_update": 0.001, "read": 0.1, "retention_per_second": 0.001, "retirement": 1.0}, "evidence_observed_at_unix_ms": now.UnixMilli(), "evidence_valid_for_ms": 600000, "horizon_seconds": 300.0, "window_cost_model": map[string]any{"implementation_id": "promql-compliance", "cost": map[string]any{"model_version": "promql-compliance-v1", "workload_fingerprint": suite.Name, "observed_at_unix_ms": now.UnixMilli(), "valid_for_ms": 600000, "horizon_seconds": 300.0, "cpu_cost": 1.0, "peak_memory_bytes": 4096, "network_bytes": 0, "storage_bytes": 2048, "source_scan_bytes": 0, "weighted_cost": 1.0}}, "scrape_interval_ms": 1000}, "environment": map[string]any{"target": "backend_local_remote_write", "collector_ids": []string{}, "capability_snapshot_id": "promql-compliance", "observed_at_unix_ms": now.UnixMilli(), "max_evidence_age_ms": 600000, "plan_version": 1, "activation_unix_ms": now.UnixMilli(), "expiry_unix_ms": nil, "backend_compat": "asap-query-backend.v1"}} +} + +func buildRevisions() (string, string, error) { + if planner, backend := os.Getenv("ASAP_PLANNER_REVISION"), os.Getenv("ASAPQUERY_BACKEND_REVISION"); planner != "" && backend != "" { + return planner, backend, nil + } + root, err := repositoryRoot() + if err != nil { + return "", "", err + } + lock, err := os.ReadFile(filepath.Join(root, "Cargo.lock")) + if err != nil { + return "", "", err + } + match := regexp.MustCompile(`(?s)name = "asap-types".*?source = "git\+[^#]+#([0-9a-f]{40})"`).FindStringSubmatch(string(lock)) + if len(match) != 2 { + return "", "", fmt.Errorf("find ASAPPlanner revision in Cargo.lock") + } + output, err := exec.Command("git", "-C", root, "rev-parse", "HEAD").Output() + if err != nil { + return "", "", fmt.Errorf("read backend revision: %w", err) + } + return match[1], strings.TrimSpace(string(output)), nil +} + +func repositoryRoot() (string, error) { + directory, err := os.Getwd() + if err != nil { + return "", err + } + for { + if _, err := os.Stat(filepath.Join(directory, "Cargo.lock")); err == nil { + return directory, nil + } + parent := filepath.Dir(directory) + if parent == directory { + return "", fmt.Errorf("find repository root from %q", directory) + } + directory = parent + } +} + +var metricToken = regexp.MustCompile(`[A-Za-z_:][A-Za-z0-9_:]*`) +var rangeToken = regexp.MustCompile(`\[([0-9]+)([smhd])\]`) + +func metricInQuery(expr string, metrics map[string]struct{}) (string, bool) { + for _, token := range metricToken.FindAllString(expr, -1) { + if _, ok := metrics[token]; ok { + return token, true + } + } + return "", false +} +func queryWindowSeconds(expr string) uint64 { + max := uint64(60) + for _, match := range rangeToken.FindAllStringSubmatch(expr, -1) { + var n uint64 + _, _ = fmt.Sscan(match[1], &n) + factor := uint64(1) + switch match[2] { + case "m": + factor = 60 + case "h": + factor = 3600 + case "d": + factor = 86400 + } + if n*factor > max { + max = n * factor + } + } + return max +} + +// PublishPhysicalPlan obtains controller manifests first, supplies unit quotes +// for the discovered components, then publishes the selected backend-local plan. +func PublishPhysicalPlan(ctx context.Context, controlPlaneURL string, request map[string]any) error { + manifests, err := postJSON(ctx, controlPlaneURL+"/api/v1/physical-plan/cost-manifests", request) + if err != nil { + return fmt.Errorf("prepare controller quotes: %w", err) + } + var entries []struct { + Components map[string]json.RawMessage `json:"components"` + } + if err := json.Unmarshal(manifests, &entries); err != nil { + return fmt.Errorf("decode controller manifests: %w", err) + } + var rawManifests []map[string]any + if err := json.Unmarshal(manifests, &rawManifests); err != nil { + return fmt.Errorf("decode controller manifests: %w", err) + } + quotes := make([]any, 0, len(entries)) + for index, entry := range entries { + costs := map[string]float64{} + for component := range entry.Components { + costs[component] = 1 + } + quotes = append(quotes, map[string]any{"manifest": rawManifests[index], "executable": index == 0, "unit_costs": costs}) + } + backendRevision, _ := request["backend_revision"].(string) + delete(request, "backend_revision") + request["workload_cost_evidence"] = map[string]any{"backend_revision": backendRevision, "planner_revision": request["planner_revision"], "data_snapshot_id": "promql-compliance", "model_version": "promql-compliance-unit-costs", "observed_at_unix_ms": time.Now().UnixMilli(), "valid_for_ms": 600_000, "quotes": quotes} + if _, err := postJSON(ctx, controlPlaneURL+"/api/v1/physical-plan/compile-and-publish", request); err != nil { + return fmt.Errorf("publish controller plan: %w", err) + } + return nil +} + +func postJSON(ctx context.Context, endpoint string, value any) ([]byte, error) { + body, err := json.Marshal(value) + if err != nil { + return nil, err + } + req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, bytes.NewReader(body)) + if err != nil { + return nil, err + } + req.Header.Set("Content-Type", "application/json") + response, err := http.DefaultClient.Do(req) + if err != nil { + return nil, err + } + defer response.Body.Close() + result, _ := io.ReadAll(io.LimitReader(response.Body, 1<<20)) + if response.StatusCode/100 != 2 { + return nil, fmt.Errorf("%s: %s", response.Status, strings.TrimSpace(string(result))) + } + return result, nil +} diff --git a/promql-compliance/runner/control_plane_test.go b/promql-compliance/runner/control_plane_test.go new file mode 100644 index 00000000..1e46b4ea --- /dev/null +++ b/promql-compliance/runner/control_plane_test.go @@ -0,0 +1,23 @@ +package runner + +import ( + "testing" + "time" +) + +func TestBuildPublicationRequestUsesSuiteQueriesAndFixtureMetric(t *testing.T) { + dataset := Dataset{Name: "fixture", Series: []DatasetSeries{{Metric: "requests_total", Samples: []DatasetSample{{OffsetSeconds: 0, Value: 1}}}}} + suite := Suite{Name: "suite", Queries: []QueryCase{{Name: "rate", Expr: "rate(requests_total[5m])", InstantOffsetsSeconds: []float64{300}}}} + request, err := BuildPublicationRequest(dataset, suite, time.Unix(1, 0)) + if err != nil { + t.Fatal(err) + } + queries := request["queries"].([]any) + query := queries[0].(map[string]any) + if query["metric"] != "requests_total" || query["window_secs"] != uint64(300) { + t.Fatalf("query = %#v", query) + } + if request["target"] != "backend_local_remote_write" || len(request["collector_ids"].([]string)) != 0 { + t.Fatalf("request = %#v", request) + } +} diff --git a/promql-compliance/runner/go.mod b/promql-compliance/runner/go.mod new file mode 100644 index 00000000..e3e5d19e --- /dev/null +++ b/promql-compliance/runner/go.mod @@ -0,0 +1,19 @@ +module github.com/ProjectASAP/ASAPQuery-backend/promql-compliance/runner + +go 1.25.8 + +require ( + github.com/gogo/protobuf v1.3.2 + github.com/golang/snappy v1.0.0 + github.com/prometheus/prometheus v0.314.0 + gopkg.in/yaml.v3 v3.0.1 +) + +require ( + github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/grafana/regexp v0.0.0-20250905093917-f7b3be9d1853 // indirect + github.com/prometheus/client_model v0.6.2 // indirect + github.com/prometheus/common v0.70.1 // indirect + golang.org/x/text v0.40.0 // indirect + google.golang.org/protobuf v1.36.12-0.20260120151049-f2248ac996af // indirect +) diff --git a/promql-compliance/runner/go.sum b/promql-compliance/runner/go.sum new file mode 100644 index 00000000..f4af4a06 --- /dev/null +++ b/promql-compliance/runner/go.sum @@ -0,0 +1,61 @@ +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM= +github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/gogo/protobuf v1.3.2 h1:Ov1cvc58UF3b5XjBnZv7+opcTcQFZebYjWzi34vdm4Q= +github.com/gogo/protobuf v1.3.2/go.mod h1:P1XiOD3dCwIKUDQYPy72D8LYyHL2YPYrpS2s69NZV8Q= +github.com/golang/snappy v1.0.0 h1:Oy607GVXHs7RtbggtPBnr2RmDArIsAefDwvrdWvRhGs= +github.com/golang/snappy v1.0.0/go.mod h1:/XxbfmMg8lxefKM7IXC3fBNl/7bRcc72aCRzEWrmP2Q= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/grafana/regexp v0.0.0-20250905093917-f7b3be9d1853 h1:cLN4IBkmkYZNnk7EAJ0BHIethd+J6LqxFNw5mSiI2bM= +github.com/grafana/regexp v0.0.0-20250905093917-f7b3be9d1853/go.mod h1:+JKpmjMGhpgPL+rXZ5nsZieVzvarn86asRlBg4uNGnk= +github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8= +github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= +github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRIccs7FGNTlIRMkT8wgtp5eCXdBlqhYGL6U= +github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk= +github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE= +github.com/prometheus/common v0.70.1 h1:1HvjP4D5oL3t8RsPlwxA9onvvStjtIHYE5XuuwOi/PY= +github.com/prometheus/common v0.70.1/go.mod h1:VdFUQDMZK3VLkurFUVhia6uys/0suUp86TJz5qbJRhc= +github.com/prometheus/prometheus v0.314.0 h1:YjsimqsIi6/mOtzZcrPEYUALO6zpfaht9O5sXqDz2vg= +github.com/prometheus/prometheus v0.314.0/go.mod h1:zjg3pMTAkY0/JG8jy/h8/YgSQUVB+aCXMhUqN6l64jg= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= +github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= +go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ= +go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ= +golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= +golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= +golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= +golang.org/x/mod v0.2.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= +golang.org/x/mod v0.3.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= +golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= +golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= +golang.org/x/net v0.0.0-20200226121028-0de0cce0169b/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= +golang.org/x/net v0.0.0-20201021035429-f5854403a974/go.mod h1:sp8m0HH+o8qH0wwXwYZr8TS3Oi6o0r6Gce1SSxlDquU= +golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.0.0-20190911185100-cd5d95a43a6e/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= +golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20200930185726-fdedc70b468f/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= +golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= +golang.org/x/text v0.40.0 h1:Ub2Z6/xjgF1WrYQz2nuITOEegKFtiIy+rieRJ5lHZKs= +golang.org/x/text v0.40.0/go.mod h1:hpnzDAfGV753zIKo+wk3u1bVKCGPbrnF7+7LBF/UHVY= +golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= +golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= +golang.org/x/tools v0.0.0-20200619180055-7c47624df98f/go.mod h1:EkVYQZoAsY45+roYkvgYkIh4xh/qjgUK9TdY2XT94GE= +golang.org/x/tools v0.0.0-20210106214847-113979e3529a/go.mod h1:emZCQorbCU4vsT4fOWvOPXz4eW1wZW4PmDk9uLelYpA= +golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= +golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= +golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= +golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= +google.golang.org/protobuf v1.36.12-0.20260120151049-f2248ac996af h1:+5/Sw3GsDNlEmu7TfklWKPdQ0Ykja5VEmq2i817+jbI= +google.golang.org/protobuf v1.36.12-0.20260120151049-f2248ac996af/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/promql-compliance/runner/lifecycle.go b/promql-compliance/runner/lifecycle.go new file mode 100644 index 00000000..58f18283 --- /dev/null +++ b/promql-compliance/runner/lifecycle.go @@ -0,0 +1,143 @@ +package runner + +import ( + "context" + "fmt" + "net/http" + "os" + "os/exec" + "path/filepath" + "strings" + "time" +) + +type ComposeLifecycle struct { + Files []string + Project string + LogsDirectory string + PlanningSnapshot string + PlanningSnapshotTemplate string + environment []string + started bool +} + +func (l *ComposeLifecycle) Start(ctx context.Context) error { + if len(l.Files) == 0 { + return nil + } + sibling, err := siblingCheckoutRoot() + if err != nil { + return err + } + environment := append(os.Environ(), + "ASAP_PRECOMPUTE_RS_CONTEXT="+filepath.Join(sibling, "ASAPCollector/asap-precompute-rs"), + "ASAP_SKETCHLIB_CONTEXT="+filepath.Join(sibling, "asap_sketchlib"), + "ASAP_GORILLA_RUST_CONTEXT="+filepath.Join(sibling, "ASAPCollector/asap-gorilla-rust"), + "ASAP_PLANNING_SNAPSHOT="+l.PlanningSnapshot, + "ASAP_PLANNING_SNAPSHOT_TEMPLATE="+l.PlanningSnapshotTemplate, + ) + l.environment = environment + // A prior interrupted run may have left data behind. Each comparison case + // must ingest into an empty Prometheus and backend state. + if err := l.runCompose(ctx, environment, "down", "--volumes", "--remove-orphans"); err != nil { + return fmt.Errorf("reset Compose services: %w", err) + } + if err := l.runCompose(ctx, environment, "up", "-d", "--build", "prometheus", "planner"); err != nil { + return err + } + l.started = true + if err := l.runCompose(ctx, environment, "wait", "planner"); err != nil { + return fmt.Errorf("derive workload cost evidence: %w", err) + } + if err := l.runCompose(ctx, environment, "up", "-d", "--build", "data-plane"); err != nil { + return err + } + return nil +} + +func (l *ComposeLifecycle) runCompose(ctx context.Context, environment []string, commandArgs ...string) error { + args := append(l.args(), commandArgs...) + command := exec.CommandContext(ctx, "docker", args...) + command.Env = environment + command.Stdout = os.Stdout + command.Stderr = os.Stderr + if err := command.Run(); err != nil { + return fmt.Errorf("start Compose: %w", err) + } + return nil +} + +func siblingCheckoutRoot() (string, error) { + root, err := repositoryRoot() + if err != nil { + return "", err + } + output, err := exec.Command("git", "-C", root, "rev-parse", "--path-format=absolute", "--git-common-dir").Output() + if err != nil { + return "", fmt.Errorf("find shared Git directory: %w", err) + } + // A linked worktree's common Git directory belongs to the primary checkout, + // whose parent is the directory containing the sibling repositories. + return filepath.Dir(filepath.Dir(strings.TrimSpace(string(output)))), nil +} +func (l *ComposeLifecycle) Stop() { + if !l.started { + return + } + l.collectLogs() + ctx, cancel := context.WithTimeout(context.Background(), time.Minute) + defer cancel() + args := append(l.args(), "down", "--volumes", "--remove-orphans") + command := exec.CommandContext(ctx, "docker", args...) + command.Env = l.environment + _ = command.Run() +} +func (l *ComposeLifecycle) collectLogs() { + if l.LogsDirectory == "" || !l.started { + return + } + _ = os.MkdirAll(l.LogsDirectory, 0o755) + command := exec.Command("docker", append(l.args(), "logs", "--no-color")...) + command.Env = l.environment + output, err := command.CombinedOutput() + if err == nil { + _ = os.WriteFile(filepath.Join(l.LogsDirectory, "compose.log"), output, 0o644) + } +} +func (l *ComposeLifecycle) args() []string { + args := []string{"compose"} + if l.Project != "" { + args = append(args, "--project-name", l.Project) + } + for _, file := range l.Files { + args = append(args, "--file", file) + } + return args +} + +func WaitForHTTP(ctx context.Context, endpoint string) error { + deadline := time.NewTimer(3 * time.Minute) + defer deadline.Stop() + ticker := time.NewTicker(time.Second) + defer ticker.Stop() + var last error + for { + response, err := http.Get(endpoint) + if err == nil { + _ = response.Body.Close() + if response.StatusCode/100 == 2 { + return nil + } + last = fmt.Errorf("%s", response.Status) + } else { + last = err + } + select { + case <-ctx.Done(): + return ctx.Err() + case <-deadline.C: + return fmt.Errorf("wait for %s: %w", endpoint, last) + case <-ticker.C: + } + } +} diff --git a/promql-compliance/runner/query.go b/promql-compliance/runner/query.go new file mode 100644 index 00000000..1e53e5d2 --- /dev/null +++ b/promql-compliance/runner/query.go @@ -0,0 +1,66 @@ +package runner + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "net/url" + "strings" + "time" +) + +// QueryResponse preserves the public Prometheus HTTP API response so the +// comparator can report target errors without discarding their payloads. +type QueryResponse struct { + Status string `json:"status"` + Data json.RawMessage `json:"data"` + ErrorType string `json:"errorType"` + Error string `json:"error"` + // ServedBy is emitted by ASAPQuery when its router answered the request. + // An empty value on the backend target means the response was forwarded to + // Prometheus, which does not emit this backend-owned header. + ServedBy string `json:"servedBy,omitempty"` +} + +type HTTPQueryTarget struct { + BaseURL string + BackendTarget bool +} + +func (target HTTPQueryTarget) Instant(ctx context.Context, expr string, at time.Time) (QueryResponse, error) { + return target.request(ctx, "/api/v1/query", url.Values{"query": {expr}, "time": {fmt.Sprintf("%.3f", float64(at.UnixMilli())/1000)}}) +} + +func (target HTTPQueryTarget) Range(ctx context.Context, expr string, spec RangeSpec, base time.Time) (QueryResponse, error) { + return target.request(ctx, "/api/v1/query_range", url.Values{ + "query": {expr}, + "start": {fmt.Sprintf("%.3f", float64(base.Add(time.Duration(spec.StartOffsetSeconds*float64(time.Second))).UnixMilli())/1000)}, + "end": {fmt.Sprintf("%.3f", float64(base.Add(time.Duration(spec.EndOffsetSeconds*float64(time.Second))).UnixMilli())/1000)}, + "step": {fmt.Sprintf("%.3f", spec.StepSeconds)}, + }) +} + +func (target HTTPQueryTarget) request(ctx context.Context, path string, query url.Values) (QueryResponse, error) { + endpoint := strings.TrimRight(target.BaseURL, "/") + path + "?" + query.Encode() + request, err := http.NewRequestWithContext(ctx, http.MethodGet, endpoint, nil) + if err != nil { + return QueryResponse{}, err + } + response, err := http.DefaultClient.Do(request) + if err != nil { + return QueryResponse{}, err + } + defer response.Body.Close() + var body QueryResponse + if err := json.NewDecoder(response.Body).Decode(&body); err != nil { + return QueryResponse{}, err + } + body.ServedBy = response.Header.Get("X-ASAP-Data-Source") + if body.ServedBy == "" && target.BackendTarget { + // Prometheus fallback returns its native response unchanged, so it has + // no ASAPQuery-owned provenance header. + body.ServedBy = "prometheus_fallback" + } + return body, nil +} diff --git a/promql-compliance/runner/remote_write.go b/promql-compliance/runner/remote_write.go new file mode 100644 index 00000000..25f0d213 --- /dev/null +++ b/promql-compliance/runner/remote_write.go @@ -0,0 +1,91 @@ +package runner + +import ( + "bytes" + "context" + "fmt" + "io" + "net/http" + "sort" + "strings" + + "github.com/gogo/protobuf/proto" + "github.com/golang/snappy" + "github.com/prometheus/prometheus/prompb" +) + +// EncodeRemoteWrite produces the canonical Prometheus Remote Write v1 body. +// Callers encode once and reuse these bytes for both comparison targets. +func EncodeRemoteWrite(baseTimeMs int64, dataset Dataset) ([]byte, error) { + request := &prompb.WriteRequest{Timeseries: make([]prompb.TimeSeries, 0, len(dataset.Series))} + for _, series := range dataset.Series { + labels := []prompb.Label{{Name: "__name__", Value: series.Metric}} + for name, value := range series.Labels { + labels = append(labels, prompb.Label{Name: name, Value: value}) + } + sort.Slice(labels, func(left, right int) bool { return labels[left].Name < labels[right].Name }) + samples := make([]prompb.Sample, 0, len(series.Samples)) + for _, sample := range series.Samples { + samples = append(samples, prompb.Sample{Value: sample.Value, Timestamp: baseTimeMs + int64(sample.OffsetSeconds*1000)}) + } + request.Timeseries = append(request.Timeseries, prompb.TimeSeries{Labels: labels, Samples: samples}) + } + encoded, err := proto.Marshal(request) + if err != nil { + return nil, fmt.Errorf("marshal Remote Write: %w", err) + } + return snappy.Encode(nil, encoded), nil +} + +func DecodeRemoteWrite(body []byte) (*prompb.WriteRequest, error) { + decoded, err := snappy.Decode(nil, body) + if err != nil { + return nil, fmt.Errorf("snappy decode Remote Write: %w", err) + } + request := &prompb.WriteRequest{} + if err := proto.Unmarshal(decoded, request); err != nil { + return nil, fmt.Errorf("unmarshal Remote Write: %w", err) + } + return request, nil +} + +// PushRemoteWrite sends an already encoded body unchanged to every target. +func PushRemoteWrite(ctx context.Context, body []byte, targets ...string) error { + for _, target := range targets { + request, err := http.NewRequestWithContext(ctx, http.MethodPost, strings.TrimRight(target, "/")+"/api/v1/write", bytes.NewReader(body)) + if err != nil { + return fmt.Errorf("build Remote Write request: %w", err) + } + request.Header.Set("Content-Type", "application/x-protobuf") + request.Header.Set("Content-Encoding", "snappy") + request.Header.Set("X-Prometheus-Remote-Write-Version", "0.1.0") + response, err := http.DefaultClient.Do(request) + if err != nil { + return fmt.Errorf("push Remote Write to %s: %w", target, err) + } + if response.StatusCode/100 != 2 { + message, _ := io.ReadAll(io.LimitReader(response.Body, 4096)) + response.Body.Close() + return fmt.Errorf("push Remote Write to %s: %s: %s", target, response.Status, message) + } + response.Body.Close() + } + return nil +} + +// Drain closes finite backend input so comparisons never race precompute work. +func Drain(ctx context.Context, target string) error { + request, err := http.NewRequestWithContext(ctx, http.MethodPost, strings.TrimRight(target, "/")+"/api/v1/precompute/drain", nil) + if err != nil { + return err + } + response, err := http.DefaultClient.Do(request) + if err != nil { + return err + } + defer response.Body.Close() + if response.StatusCode/100 != 2 { + return fmt.Errorf("drain %s: %s", target, response.Status) + } + return nil +} diff --git a/promql-compliance/runner/report.go b/promql-compliance/runner/report.go new file mode 100644 index 00000000..9c536b24 --- /dev/null +++ b/promql-compliance/runner/report.go @@ -0,0 +1,121 @@ +package runner + +import ( + "context" + "fmt" + "time" +) + +// Report is written even when query comparisons fail so a manual CI run has +// evidence for every attempted evaluation. +type Report struct { + Suite string `json:"suite"` + Dataset string `json:"dataset"` + BaseTime time.Time `json:"baseTime"` + Queries []QueryReport `json:"queries"` + Passed bool `json:"passed"` +} + +type QueryReport struct { + Name string `json:"name"` + Expr string `json:"expr"` + Tolerance ComparisonPolicy `json:"tolerance"` + Range *ComparisonOutcome `json:"range,omitempty"` + RangeResponses *ResponsePair `json:"rangeResponses,omitempty"` + Instant []InstantComparison `json:"instant,omitempty"` + ReferenceParity []InstantComparison `json:"referenceParity,omitempty"` + BackendParity []InstantComparison `json:"backendParity,omitempty"` + Passed bool `json:"passed"` +} + +type InstantComparison struct { + OffsetSeconds float64 `json:"offsetSeconds"` + Time time.Time `json:"time"` + Comparison ComparisonOutcome `json:"comparison"` + Responses ResponsePair `json:"responses"` +} + +// ResponsePair preserves the raw public API payloads behind each comparison. +// It lets a passing report show exactly what Prometheus and ASAPQuery returned. +type ResponsePair struct { + Reference QueryResponse `json:"reference"` + Backend QueryResponse `json:"backend"` +} +type ComparisonOutcome struct { + Passed bool `json:"passed"` + Diff string `json:"diff,omitempty"` + ReferenceError string `json:"referenceError,omitempty"` + BackendError string `json:"backendError,omitempty"` +} + +func CompareSuite(ctx context.Context, reference, backend HTTPQueryTarget, suite Suite, dataset string, base time.Time) Report { + report := Report{Suite: suite.Name, Dataset: dataset, BaseTime: base, Passed: true} + for _, query := range suite.Queries { + item := CompareQuery(ctx, reference, backend, query, suite.ComparisonDefaults, base) + report.Queries = append(report.Queries, item) + report.Passed = report.Passed && item.Passed + } + return report +} + +func CompareQuery(ctx context.Context, reference, backend HTTPQueryTarget, query QueryCase, defaults ComparisonPolicy, base time.Time) QueryReport { + policy := query.EffectiveTolerance(defaults) + report := QueryReport{Name: query.Name, Expr: query.Expr, Tolerance: policy, Passed: true} + var referenceRange, backendRange QueryResponse + if query.Range != nil { + left, leftErr := reference.Range(ctx, query.Expr, *query.Range, base) + right, rightErr := backend.Range(ctx, query.Expr, *query.Range, base) + outcome := compareOutcome(left, right, leftErr, rightErr, policy) + report.Range, report.Passed = &outcome, report.Passed && outcome.Passed + report.RangeResponses = &ResponsePair{Reference: left, Backend: right} + referenceRange, backendRange = left, right + } + for _, offset := range query.InstantOffsetsSeconds { + at := base.Add(time.Duration(offset * float64(time.Second))) + left, leftErr := reference.Instant(ctx, query.Expr, at) + right, rightErr := backend.Instant(ctx, query.Expr, at) + outcome := compareOutcome(left, right, leftErr, rightErr, policy) + report.Instant = append(report.Instant, InstantComparison{OffsetSeconds: offset, Time: at, Comparison: outcome, Responses: ResponsePair{Reference: left, Backend: right}}) + report.Passed = report.Passed && outcome.Passed + if query.Range == nil || leftErr != nil || rightErr != nil || referenceRange.Status != "success" || backendRange.Status != "success" || left.Status != "success" || right.Status != "success" { + continue + } + refParity := parityOutcome(referenceRange, left, at, policy) + backendParity := parityOutcome(backendRange, right, at, policy) + report.ReferenceParity = append(report.ReferenceParity, InstantComparison{OffsetSeconds: offset, Time: at, Comparison: refParity}) + report.BackendParity = append(report.BackendParity, InstantComparison{OffsetSeconds: offset, Time: at, Comparison: backendParity}) + report.Passed = report.Passed && refParity.Passed && backendParity.Passed + } + return report +} + +func compareOutcome(left, right QueryResponse, leftErr, rightErr error, policy ComparisonPolicy) ComparisonOutcome { + outcome := ComparisonOutcome{} + if leftErr != nil { + outcome.ReferenceError = leftErr.Error() + } + if rightErr != nil { + outcome.BackendError = rightErr.Error() + } + if leftErr == nil && rightErr == nil { + if err := CompareResponses(left, right, policy); err != nil { + outcome.Diff = err.Error() + } + if outcome.Diff == "" && right.ServedBy == "prometheus_fallback" { + outcome.Diff = "backend forwarded the query to Prometheus fallback; no ASAPQuery result was compared" + } + if outcome.Diff == "" && right.ServedBy == "" { + outcome.Diff = "backend response has no serving provenance" + } + } + outcome.Passed = outcome.Diff == "" && outcome.ReferenceError == "" && outcome.BackendError == "" + return outcome +} + +func parityOutcome(rangeResponse, instantResponse QueryResponse, at time.Time, policy ComparisonPolicy) ComparisonOutcome { + err := CompareRangeAtInstant(rangeResponse, instantResponse, at, policy) + if err == nil { + return ComparisonOutcome{Passed: true} + } + return ComparisonOutcome{Diff: fmt.Sprintf("range/instant parity: %v", err)} +} diff --git a/promql-compliance/runner/report_test.go b/promql-compliance/runner/report_test.go new file mode 100644 index 00000000..474c47e8 --- /dev/null +++ b/promql-compliance/runner/report_test.go @@ -0,0 +1,12 @@ +package runner + +import "testing" + +func TestCompareOutcomeRejectsPrometheusFallback(t *testing.T) { + response := comparisonResponse(t, "vector", `[]`) + response.ServedBy = "prometheus_fallback" + outcome := compareOutcome(comparisonResponse(t, "vector", `[]`), response, nil, nil, ComparisonPolicy{}) + if outcome.Passed || outcome.Diff == "" { + t.Fatalf("fallback must not count as an ASAPQuery comparison: %+v", outcome) + } +} diff --git a/promql-compliance/runner/spec.go b/promql-compliance/runner/spec.go new file mode 100644 index 00000000..84b2ed86 --- /dev/null +++ b/promql-compliance/runner/spec.go @@ -0,0 +1,216 @@ +// Package runner owns the declarative inputs for backend PromQL differential +// tests. The runner will derive both control-plane publication input and query +// evaluations from this single suite definition. +package runner + +import ( + "bytes" + "fmt" + "math" + "os" + "sort" + "strings" + + "gopkg.in/yaml.v3" +) + +type Suite struct { + Name string `yaml:"name"` + ComparisonDefaults ComparisonPolicy `yaml:"comparison_defaults"` + Queries []QueryCase `yaml:"queries"` +} + +// Dataset contains source samples expressed relative to the run's base time. +// The same expanded Remote Write payload will later be sent to both targets. +type Dataset struct { + Name string `yaml:"name"` + Series []DatasetSeries `yaml:"series"` +} + +type DatasetSeries struct { + Metric string `yaml:"metric"` + Labels map[string]string `yaml:"labels"` + Samples []DatasetSample `yaml:"samples"` +} + +type DatasetSample struct { + OffsetSeconds float64 `yaml:"offset_seconds"` + Value float64 `yaml:"value"` +} + +type QueryCase struct { + Name string `yaml:"name"` + Expr string `yaml:"expr"` + InstantOffsetsSeconds []float64 `yaml:"instant_offsets_seconds"` + Range *RangeSpec `yaml:"range"` + Comparison *ComparisonPolicy `yaml:"comparison"` +} + +type RangeSpec struct { + StartOffsetSeconds float64 `yaml:"start_offset_seconds"` + EndOffsetSeconds float64 `yaml:"end_offset_seconds"` + StepSeconds float64 `yaml:"step_seconds"` +} + +type ComparisonPolicy struct { + ValueTolerance *Tolerance `yaml:"value_tolerance"` +} + +type Tolerance struct { + Relative *float64 `yaml:"relative"` + Absolute *float64 `yaml:"absolute"` +} + +// LoadSuite rejects implicit evaluation windows. Fixed evaluation timestamps +// are required so that Prometheus and the backend see the same query input. +func LoadSuite(contents []byte) (Suite, error) { + var suite Suite + decoder := yaml.NewDecoder(bytes.NewReader(contents)) + decoder.KnownFields(true) + if err := decoder.Decode(&suite); err != nil { + return Suite{}, fmt.Errorf("parse query suite: %w", err) + } + if suite.Name == "" { + return Suite{}, fmt.Errorf("query suite has no name") + } + if len(suite.Queries) == 0 { + return Suite{}, fmt.Errorf("query suite %q has no queries", suite.Name) + } + for index := range suite.Queries { + query := &suite.Queries[index] + if query.Name == "" || query.Expr == "" { + return Suite{}, fmt.Errorf("query %d requires name and expr", index) + } + if len(query.InstantOffsetsSeconds) == 0 && query.Range == nil { + return Suite{}, fmt.Errorf("query %q has neither instant times nor a range", query.Name) + } + if query.Range != nil { + if err := query.Range.validate(); err != nil { + return Suite{}, fmt.Errorf("query %q: %w", query.Name, err) + } + for _, offset := range query.InstantOffsetsSeconds { + if !finite(offset) || offset < query.Range.StartOffsetSeconds || offset > query.Range.EndOffsetSeconds { + return Suite{}, fmt.Errorf("query %q has instant offset outside its range", query.Name) + } + } + } else { + for _, offset := range query.InstantOffsetsSeconds { + if !finite(offset) { + return Suite{}, fmt.Errorf("query %q has non-finite instant offset", query.Name) + } + } + } + if err := validateTolerance(query.EffectiveTolerance(suite.ComparisonDefaults).ValueTolerance); err != nil { + return Suite{}, fmt.Errorf("query %q: %w", query.Name, err) + } + } + return suite, nil +} + +// LoadDataset validates the properties needed to construct deterministic +// Remote Write batches: one unique label set per metric and strictly ordered, +// finite samples per series. +func LoadDataset(contents []byte) (Dataset, error) { + var dataset Dataset + decoder := yaml.NewDecoder(bytes.NewReader(contents)) + decoder.KnownFields(true) + if err := decoder.Decode(&dataset); err != nil { + return Dataset{}, fmt.Errorf("parse dataset: %w", err) + } + if dataset.Name == "" || len(dataset.Series) == 0 { + return Dataset{}, fmt.Errorf("dataset requires name and series") + } + seen := make(map[string]struct{}, len(dataset.Series)) + for index, series := range dataset.Series { + if series.Metric == "" || len(series.Samples) == 0 { + return Dataset{}, fmt.Errorf("series %d requires metric and samples", index) + } + key := series.Metric + "\x00" + canonicalLabels(series.Labels) + if _, duplicate := seen[key]; duplicate { + return Dataset{}, fmt.Errorf("dataset has duplicate series %q", series.Metric) + } + seen[key] = struct{}{} + var previous float64 + for sampleIndex, sample := range series.Samples { + if !finite(sample.OffsetSeconds) || !finite(sample.Value) { + return Dataset{}, fmt.Errorf("series %q sample %d is non-finite", series.Metric, sampleIndex) + } + if sampleIndex > 0 && sample.OffsetSeconds <= previous { + return Dataset{}, fmt.Errorf("series %q samples are not strictly ordered", series.Metric) + } + previous = sample.OffsetSeconds + } + } + return dataset, nil +} + +func LoadSuiteFile(path string) (Suite, error) { + contents, err := os.ReadFile(path) + if err != nil { + return Suite{}, err + } + return LoadSuite(contents) +} +func LoadDatasetFile(path string) (Dataset, error) { + contents, err := os.ReadFile(path) + if err != nil { + return Dataset{}, err + } + return LoadDataset(contents) +} + +func (q QueryCase) EffectiveTolerance(defaults ComparisonPolicy) ComparisonPolicy { + if q.Comparison == nil || q.Comparison.ValueTolerance == nil { + return defaults + } + result := defaults + if result.ValueTolerance == nil { + result.ValueTolerance = &Tolerance{} + } + merged := *result.ValueTolerance + if q.Comparison.ValueTolerance.Relative != nil { + merged.Relative = q.Comparison.ValueTolerance.Relative + } + if q.Comparison.ValueTolerance.Absolute != nil { + merged.Absolute = q.Comparison.ValueTolerance.Absolute + } + result.ValueTolerance = &merged + return result +} + +func (r RangeSpec) validate() error { + if !finite(r.StartOffsetSeconds) || !finite(r.EndOffsetSeconds) || !finite(r.StepSeconds) { + return fmt.Errorf("range offsets and step must be finite") + } + if r.EndOffsetSeconds <= r.StartOffsetSeconds || r.StepSeconds <= 0 { + return fmt.Errorf("range end must be after start and step must be positive") + } + return nil +} + +func validateTolerance(tolerance *Tolerance) error { + if tolerance == nil { + return nil + } + for _, value := range []*float64{tolerance.Relative, tolerance.Absolute} { + if value != nil && (!finite(*value) || *value < 0) { + return fmt.Errorf("tolerance must be finite and non-negative") + } + } + return nil +} + +func finite(value float64) bool { return !math.IsNaN(value) && !math.IsInf(value, 0) } + +func canonicalLabels(labels map[string]string) string { + keys := make([]string, 0, len(labels)) + for key := range labels { + keys = append(keys, key) + } + sort.Strings(keys) + parts := make([]string, 0, len(keys)) + for _, key := range keys { + parts = append(parts, key+"="+labels[key]) + } + return strings.Join(parts, "\x00") +} diff --git a/promql-compliance/runner/spec_test.go b/promql-compliance/runner/spec_test.go new file mode 100644 index 00000000..99c566f0 --- /dev/null +++ b/promql-compliance/runner/spec_test.go @@ -0,0 +1,114 @@ +package runner + +import "testing" + +func TestCheckedInFixturesMeetStrictContracts(t *testing.T) { + for _, path := range []string{ + "../datasets/single-rate.yaml", "../datasets/sparse-checkout.yaml", + "../datasets/aggregations.yaml", "../datasets/aggregations-dense-cadence.yaml", + } { + if _, err := LoadDatasetFile(path); err != nil { + t.Fatalf("LoadDatasetFile(%q): %v", path, err) + } + } + for _, path := range []string{"../suites/temporal.yaml", "../suites/aggregations.yaml"} { + if _, err := LoadSuiteFile(path); err != nil { + t.Fatalf("LoadSuiteFile(%q): %v", path, err) + } + } +} + +func TestLoadSuiteRequiresExplicitEvaluationAndCarriesTolerance(t *testing.T) { + suite, err := LoadSuite([]byte(`name: temporal +comparison_defaults: + value_tolerance: + relative: 0.01 +queries: + - name: request-rate + expr: rate(http_requests_total[5m]) + instant_offsets_seconds: [300, 600] + range: + start_offset_seconds: 300 + end_offset_seconds: 600 + step_seconds: 60 +`)) + if err != nil { + t.Fatalf("LoadSuite: %v", err) + } + if got, want := len(suite.Queries), 1; got != want { + t.Fatalf("query count = %d, want %d", got, want) + } + if got := suite.Queries[0].EffectiveTolerance(suite.ComparisonDefaults); got.ValueTolerance == nil || got.ValueTolerance.Relative == nil || *got.ValueTolerance.Relative != 0.01 { + t.Fatalf("effective tolerance = %#v, want inherited relative tolerance", got) + } +} + +func TestLoadSuiteRejectsImplicitEvaluation(t *testing.T) { + _, err := LoadSuite([]byte(`name: incomplete +queries: + - name: no-evaluation + expr: up +`)) + if err == nil { + t.Fatal("LoadSuite accepted a query without an instant time or range") + } +} + +func TestLoadSuiteRejectsRemovedExpectErrorField(t *testing.T) { + _, err := LoadSuite([]byte(`name: invalid +queries: + - name: request-rate + expr: rate(http_requests_total[5m]) + instant_offsets_seconds: [300] + expect_error: true +`)) + if err == nil { + t.Fatal("LoadSuite accepted the removed expect_error field") + } +} + +func TestLoadDatasetRejectsDuplicateSeriesAndUnorderedSamples(t *testing.T) { + _, err := LoadDataset([]byte(`name: invalid +series: + - metric: requests_total + labels: {host: a} + samples: + - {offset_seconds: 60, value: 1} + - {offset_seconds: 0, value: 0} + - metric: requests_total + labels: {host: a} + samples: + - {offset_seconds: 0, value: 0} +`)) + if err == nil { + t.Fatal("LoadDataset accepted unordered or duplicate series") + } +} + +func TestEncodeRemoteWriteRoundTripsFixtureSamples(t *testing.T) { + dataset, err := LoadDataset([]byte(`name: request +series: + - metric: requests_total + labels: {host: a} + samples: + - {offset_seconds: 0, value: 1} + - {offset_seconds: 60, value: 2} +`)) + if err != nil { + t.Fatal(err) + } + body, err := EncodeRemoteWrite(1_700_000_000_000, dataset) + if err != nil { + t.Fatal(err) + } + decoded, err := DecodeRemoteWrite(body) + if err != nil { + t.Fatal(err) + } + if got, want := len(decoded.Timeseries), 1; got != want { + t.Fatalf("series = %d, want %d", got, want) + } + if got, want := decoded.Timeseries[0].Samples[1].Timestamp, int64(1_700_000_060_000); got != want { + t.Fatalf("timestamp = %d, want %d", got, want) + } +} diff --git a/promql-compliance/suites/aggregations.yaml b/promql-compliance/suites/aggregations.yaml new file mode 100644 index 00000000..af11f833 --- /dev/null +++ b/promql-compliance/suites/aggregations.yaml @@ -0,0 +1,288 @@ +name: aggregations +comparison_defaults: + value_tolerance: + relative: 0.01 + absolute: 0.000001 + +queries: + # Every query is exercised as both instant and range. Keep the evaluation + # grid shared so the runner's range/instant parity check applies uniformly. + - name: increase + expr: "increase(data[5m])" + instant_offsets_seconds: &evaluation_offsets [300, 600, 900, 1200] + range: &evaluation_range + start_offset_seconds: 300 + end_offset_seconds: 1200 + step_seconds: 60 + + - name: rate + expr: "rate(data[5m])" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + + - name: quantile-over-time-0.1 + expr: "quantile_over_time(0.1, data[5m])" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-over-time-0.2 + expr: "quantile_over_time(0.2, data[5m])" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-over-time-0.3 + expr: "quantile_over_time(0.3, data[5m])" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-over-time-0.4 + expr: "quantile_over_time(0.4, data[5m])" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-over-time-0.5 + expr: "quantile_over_time(0.5, data[5m])" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-over-time-0.6 + expr: "quantile_over_time(0.6, data[5m])" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-over-time-0.7 + expr: "quantile_over_time(0.7, data[5m])" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-over-time-0.8 + expr: "quantile_over_time(0.8, data[5m])" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-over-time-0.9 + expr: "quantile_over_time(0.9, data[5m])" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-over-time-0.95 + expr: "quantile_over_time(0.95, data[5m])" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-over-time-0.99 + expr: "quantile_over_time(0.99, data[5m])" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + + - name: sum-over-time + expr: "sum_over_time(data[5m])" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: count-over-time + expr: "count_over_time(data[5m])" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + + # Bare-selector topk (no wrapped temporal function): unlike every + # topk-over-temporal case below, this keeps __name__ -- see #710. + - name: topk-1-bare + expr: "topk(1, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-3-bare + expr: "topk(3, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-5-bare + expr: "topk(5, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-10-bare + expr: "topk(10, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + + # Zero-label grouping is the ordinary aggregate form. One- and two-label + # forms retain job, then job+instance, respectively. + - name: topk-1-sum-over-time + expr: "topk(1, sum_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-3-sum-over-time + expr: "topk(3, sum_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-5-sum-over-time + expr: "topk(5, sum_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-10-sum-over-time + expr: "topk(10, sum_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-1-sum-over-time-by-job + expr: "topk by (job) (1, sum_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-3-sum-over-time-by-job + expr: "topk by (job) (3, sum_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-5-sum-over-time-by-job + expr: "topk by (job) (5, sum_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-10-sum-over-time-by-job + expr: "topk by (job) (10, sum_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-1-sum-over-time-by-job-instance + expr: "topk by (job, instance) (1, sum_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-3-sum-over-time-by-job-instance + expr: "topk by (job, instance) (3, sum_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-5-sum-over-time-by-job-instance + expr: "topk by (job, instance) (5, sum_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-10-sum-over-time-by-job-instance + expr: "topk by (job, instance) (10, sum_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + + - name: topk-1-count-over-time + expr: "topk(1, count_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-3-count-over-time + expr: "topk(3, count_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-5-count-over-time + expr: "topk(5, count_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-10-count-over-time + expr: "topk(10, count_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-1-count-over-time-by-job + expr: "topk by (job) (1, count_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-3-count-over-time-by-job + expr: "topk by (job) (3, count_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-5-count-over-time-by-job + expr: "topk by (job) (5, count_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-10-count-over-time-by-job + expr: "topk by (job) (10, count_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-1-count-over-time-by-job-instance + expr: "topk by (job, instance) (1, count_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-3-count-over-time-by-job-instance + expr: "topk by (job, instance) (3, count_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-5-count-over-time-by-job-instance + expr: "topk by (job, instance) (5, count_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: topk-10-count-over-time-by-job-instance + expr: "topk by (job, instance) (10, count_over_time(data[5m]))" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + + - name: quantile-0.1 + expr: "quantile(0.1, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.2 + expr: "quantile(0.2, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.3 + expr: "quantile(0.3, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.4 + expr: "quantile(0.4, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.5 + expr: "quantile(0.5, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.6 + expr: "quantile(0.6, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.7 + expr: "quantile(0.7, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.8 + expr: "quantile(0.8, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.9 + expr: "quantile(0.9, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.95 + expr: "quantile(0.95, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.99 + expr: "quantile(0.99, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.5-by-job + expr: "quantile by (job) (0.5, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.95-by-job + expr: "quantile by (job) (0.95, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.99-by-job + expr: "quantile by (job) (0.99, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.5-by-job-instance + expr: "quantile by (job, instance) (0.5, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.95-by-job-instance + expr: "quantile by (job, instance) (0.95, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: quantile-0.99-by-job-instance + expr: "quantile by (job, instance) (0.99, data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + + - name: sum + expr: "sum(data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: sum-by-job + expr: "sum by (job) (data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: sum-by-job-instance + expr: "sum by (job, instance) (data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: count + expr: "count(data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: count-by-job + expr: "count by (job) (data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range + - name: count-by-job-instance + expr: "count by (job, instance) (data)" + instant_offsets_seconds: *evaluation_offsets + range: *evaluation_range diff --git a/promql-compliance/suites/temporal.yaml b/promql-compliance/suites/temporal.yaml new file mode 100644 index 00000000..69740cfc --- /dev/null +++ b/promql-compliance/suites/temporal.yaml @@ -0,0 +1,13 @@ +name: temporal +comparison_defaults: + value_tolerance: + relative: 0 + absolute: 0.000001 +queries: + - name: request-rate + expr: rate(http_requests_total[5m]) + instant_offsets_seconds: [300, 600, 1200] + range: + start_offset_seconds: 300 + end_offset_seconds: 1200 + step_seconds: 60