diff --git a/.github/workflows/promql-compliance.yml b/.github/workflows/promql-compliance.yml new file mode 100644 index 000000000..2e645f9d9 --- /dev/null +++ b/.github/workflows/promql-compliance.yml @@ -0,0 +1,45 @@ +name: PromQL compliance + +on: + pull_request: + types: [opened, synchronize, reopened, ready_for_review] + 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: dtolnay/rust-toolchain@stable + - name: Install protobuf compiler + run: sudo apt-get update && sudo apt-get install -y protobuf-compiler + - name: Test Rust compliance harness + working-directory: ASAPQuery-backend + run: cargo test --locked -p promql-compliance + - 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: Execute every admitted individual and ensemble candidate + working-directory: ASAPQuery-backend/promql-compliance/runner + run: make candidates 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/Cargo.lock b/Cargo.lock index 5e0173958..16b6036db 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3054,6 +3054,26 @@ dependencies = [ "thiserror 1.0.69", ] +[[package]] +name = "promql-compliance" +version = "0.1.0" +dependencies = [ + "anyhow", + "asap-types", + "asap_types", + "clap", + "control_plane", + "prost", + "regex", + "reqwest", + "serde", + "serde_json", + "serde_yaml", + "snap", + "tempfile", + "tokio", +] + [[package]] name = "promql-parser" version = "0.10.0" diff --git a/Cargo.toml b/Cargo.toml index 9040dbf29..cb8c2d601 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -5,6 +5,7 @@ members = [ "crates/asap_types", "data_plane", "control_plane", + "promql-compliance/runner", ] [workspace.package] diff --git a/control_plane/Dockerfile b/control_plane/Dockerfile index 0ca9f984f..ba73508be 100644 --- a/control_plane/Dockerfile +++ b/control_plane/Dockerfile @@ -15,8 +15,7 @@ RUN apt-get update && \ rm -rf /var/lib/apt/lists/* COPY . ASAPQuery-backend - -RUN cd ASAPQuery-backend && cargo build --locked --release --bin control_plane +RUN cd ASAPQuery-backend && cargo build --locked --release --bin control_plane --bin control_plane_quote_snapshot --example compile_workload_artifact # Runtime image: binary + CA certs. FROM debian:bookworm-slim @@ -25,6 +24,10 @@ 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 +COPY --from=build /src/ASAPQuery-backend/target/release/examples/compile_workload_artifact \ + /usr/local/bin/compile_workload_artifact 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 000000000..adac02476 --- /dev/null +++ b/control_plane/src/bin/control_plane_quote_snapshot.rs @@ -0,0 +1,111 @@ +//! 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, DeploymentPlanCompiler, 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 = DeploymentPlanCompiler + .compile_promql(candidate.clone(), environment.clone()) + .ok()?; + // This execution suite must exercise maintained candidates when legal. + let unit_cost = if plan.precompute_plan.materializations.is_empty() { + 1e12 + } else { + 1.0 + }; + let manifest = manifest(&plan, &candidate.queries).ok()?; + Some(WorkloadQuote { + unit_costs: manifest + .components + .keys() + .map(|key| (key.clone(), unit_cost)) + .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 + .as_ref() + .unwrap() + .quotes + .is_empty()); + let plan = quoted.compile_promql().unwrap(); + assert!( + !plan.precompute_plan.materializations.is_empty(), + "deterministic compliance quotes should select a maintained candidate" + ); + } +} diff --git a/data_plane/Dockerfile b/data_plane/Dockerfile index 004c7947a..514708e02 100644 --- a/data_plane/Dockerfile +++ b/data_plane/Dockerfile @@ -20,8 +20,11 @@ RUN apt-get update && \ rm -rf /var/lib/apt/lists/* COPY . ASAPQuery-backend - -RUN cd ASAPQuery-backend && cargo build --locked --release --bin data_plane +RUN --mount=type=cache,target=/usr/local/cargo/registry \ + --mount=type=cache,target=/usr/local/cargo/git \ + --mount=type=cache,target=/src/ASAPQuery-backend/target \ + cd ASAPQuery-backend && cargo build --locked --release --bin data_plane && \ + install -m 755 target/release/data_plane /usr/local/bin/data_plane # Runtime image: binary + CA certs. FROM debian:bookworm-slim @@ -30,7 +33,7 @@ RUN apt-get update && \ rm -rf /var/lib/apt/lists/* # Default --output-dir for the rolling log file. RUN mkdir -p /var/log/asap -COPY --from=build /src/ASAPQuery-backend/target/release/data_plane \ +COPY --from=build /usr/local/bin/data_plane \ /usr/local/bin/data_plane ENV RUST_LOG=info \ diff --git a/data_plane/src/main.rs b/data_plane/src/main.rs index 7fcefac33..f40b8f1d9 100644 --- a/data_plane/src/main.rs +++ b/data_plane/src/main.rs @@ -461,11 +461,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 { @@ -1392,6 +1387,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" + ); + assert!( Args::try_parse_from(["data_plane", "--streaming-config", "streaming.yaml"]).is_err() ); diff --git a/data_plane/tests/asapquery_compatibility_process_e2e.rs b/data_plane/tests/asapquery_compatibility_process_e2e.rs index b3cd98c7e..5555c5bb1 100644 --- a/data_plane/tests/asapquery_compatibility_process_e2e.rs +++ b/data_plane/tests/asapquery_compatibility_process_e2e.rs @@ -210,7 +210,10 @@ fn unused_port() -> u16 { } async fn wait_until_ready(client: &reqwest::Client, url: &str, child: &mut Child) { - for _ in 0..120 { + // Startup compiles the complete workload candidate set before serving. + // This is a readiness budget, not a query latency assertion. + let deadline = tokio::time::Instant::now() + Duration::from_secs(120); + while tokio::time::Instant::now() < deadline { if let Some(status) = child.try_wait().expect("inspect backend process") { panic!("backend exited before readiness: {status}"); } diff --git a/data_plane/tests/support/issue_701_702_process.rs b/data_plane/tests/support/issue_701_702_process.rs index f2dbe6b6c..d6fc4a31d 100644 --- a/data_plane/tests/support/issue_701_702_process.rs +++ b/data_plane/tests/support/issue_701_702_process.rs @@ -346,10 +346,21 @@ fn issue_701_702_uncertified_ratios_require_exact_fallback() { let candidates = workload_cost::enumerate_exact_and_materialized_candidates(request).unwrap(); assert!(!candidates.is_empty()); + let mut exact_count = 0; for candidate in candidates { - let plan = DeploymentPlanCompiler - .compile_promql(candidate, environment.clone()) - .unwrap(); + let plan = match DeploymentPlanCompiler.compile_promql(candidate, environment.clone()) { + Ok(plan) => plan, + Err(error) => { + assert!( + error + .to_string() + .contains("no certified accuracy guarantee"), + "{query}: {error}" + ); + continue; + } + }; + exact_count += 1; assert!(plan.precompute_plan.materializations.is_empty(), "{query}"); assert!( plan.query_plan.entries.values().all(|entry| matches!( @@ -359,6 +370,10 @@ fn issue_701_702_uncertified_ratios_require_exact_fallback() { "{query}" ); } + assert!( + exact_count > 0, + "{query}: exact execution must remain available" + ); } } diff --git a/docs/design_docs/evidence-dependent-candidates.md b/docs/design_docs/evidence-dependent-candidates.md index d4baf36cc..1a07d4995 100644 --- a/docs/design_docs/evidence-dependent-candidates.md +++ b/docs/design_docs/evidence-dependent-candidates.md @@ -17,14 +17,16 @@ could prove valid or deploys one whose guarantee has never been established. | --- | --- | --- | | Construct semantic candidates | Planner | Preserve unknown guarantees; reject known-invalid evidence and impossible shapes | | Supply external facts | Backend | Bind evidence to the query, data population, snapshot and validity period | -| Derive accuracy and select logical roots | Planner, under backend models and policy | Respect the root accuracy target; missing proof is not certification | +| Derive accuracy and expose supported candidates | Planner, under supplied facts and policy | Respect the root accuracy target; missing proof is not certification | | Bind and admit a deployment | Backend | Verify concrete execution support and compute complete workload resource cost | -The backend still invokes Planner's workload search and global selection. It -does not introduce a second semantic optimizer. Physical binding preserves the -selected DAG's operators, grouping, windows and dependencies; a semantic change -requires a new selection. Single-query and workload selection use the same -costed Planner search; the first-candidate helper has been removed. +Planner supplies the supported computation/physical candidate inventory. Backend +binds and prices candidates, then selects a deployment. Physical binding preserves +operators, grouping, windows and dependencies; it does not introduce another +semantic optimizer. The current validation milestone uses synthetic complete +quotes followed by execution of the selected typed plan. Online ERP observations, +feedback-driven replanning and deployment replacement are deferred. Optional +existing ERP adapters described below are not prerequisites for this milestone. ## Decision flow @@ -34,7 +36,7 @@ flowchart TD Validate -->|Invalid supplied evidence| Error[Reject request with reason] Validate -->|Valid or absent evidence| Search[Planner search retains constructible candidates] Search --> Inspect[Explain known and unknown candidate properties] - Search --> Select[Planner global selection under backend policy] + Search --> Select[Backend candidate binding and cost selection] Select --> Exact[Explicit exact fallback when no certified summary is selected] Select --> Bind[Bind selected logical DAG to concrete execution] Bind --> Admit[Check guarantees, runtime support and computed workload cost] diff --git a/docs/design_docs/physical-operators.md b/docs/design_docs/physical-operators.md index 2e47b59ad..64b3b0c1f 100644 --- a/docs/design_docs/physical-operators.md +++ b/docs/design_docs/physical-operators.md @@ -68,6 +68,27 @@ Every installed range evaluation uses the shared DAG execution path and reports its actual local-summary and external-exact work. The differential and benefit runners require successful local execution evidence as well as matching values. Routing to the ASAP endpoint alone is insufficient. + +## Planner physical candidates and precompute boundaries + +Planner #462 exposes `physical_planner::compile_candidates` and +`select_candidate`. A candidate carries precompute/query Physical DAGs and +typed materialized outputs. For grouped rate, Planner can compile both +per-series Rate → stored values → query-side Sum and precompute Rate → Sum → +stored grouped values. Counter-state inputs remain explicit maintenance +dependencies; raw-counter Sum before Rate is not equivalent. + +The backend reports source/state/operator feasibility and complete workload +costs. `implementation.require_backend_local_execution` rejects external exact +dependencies before pricing when the deployment has no such service. The runner +sets this requirement because its deployment disables forwarding, and declares +actual replay rate/cadence rather than synthetic compatibility demand. + +Planner capability does not prove that the backend can persist every value-output +frontier. Result-row publication, revision/coverage and retention bindings must +be admitted explicitly; the deployment compiler must reject unsupported +frontiers instead of moving operators. + Selected Filter, Sort, Limit and semi-join fragments are compiled by Planner and persisted with typed input contracts. Complete query candidates include Planner readouts that turn accumulator state into values; internal shared and stored edges retain diff --git a/docs/design_docs/promql-compliance-grill.md b/docs/design_docs/promql-compliance-grill.md new file mode 100644 index 000000000..42881d029 --- /dev/null +++ b/docs/design_docs/promql-compliance-grill.md @@ -0,0 +1,87 @@ +# Current execution scope + +The current milestone is workload → Planner physical candidates → Backend selection +with synthetic costs → typed installation → data-plane execution. The original +Q&A below records earlier decisions. Online ERP collection, query-time statistics +feedback and runtime replanning are deferred. Synthetic prices never weaken the +accuracy, local-execution or bound-SDS contracts. + +# 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:** how should the fixture seeder, comparator and Docker Compose runner integrate with the Rust backend? +> +> **Decision:** implement the runner in the Cargo workspace. Shared Rust modules provide deterministic Remote Write encoding, structured comparison, fixtures, tolerances, reports and Compose lifecycle handling. No Go toolchain is required. + +**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/docs/evaluation/candidate-ensembles-2026-09-28/README.md b/docs/evaluation/candidate-ensembles-2026-09-28/README.md new file mode 100644 index 000000000..c840c6fb8 --- /dev/null +++ b/docs/evaluation/candidate-ensembles-2026-09-28/README.md @@ -0,0 +1,53 @@ +# Individual and ensemble candidate execution + +Every admitted executable candidate passed selection, typed deployment installation, +Remote Write ingestion, precompute drain and comparison with Prometheus. This is +fixture execution evidence, not production cost validation or human plan approval. + +| Workload | Candidate runs | +| --- | --- | +| query:spatial-sum | 1 | +| query:spatial-topk | 1 | +| query:spatial-quantile | 1 | +| query:temporal-sum | 1 | +| query:temporal-quantile | 1 | +| query:temporal-rate | 1 | +| query:grouped-rate | 3 | +| query:grouped-temporal-sum | 2 | +| query:topk-rate | 2 | +| query:quantile-ratio | 1 | +| shared-rate | 2 | +| shared-quantiles | 1 | +| full-ensemble | 3 | + +Total: 20 candidate runs; 52 query comparisons. + +For each target, only synthetic prices change. Production selection must return +its exact manifest. An ensemble is installed once and ingested once per candidate; +all member queries run against that installation. Generated admission reports +include rejected candidates. Coverage is the exposed, admitted fixture inventory, not an +exhaustive Cartesian product of theoretical physical plans. + +Validation also passed: three Level 1 tests, the complete Level 2 ranking test, +14 harness contract tests, three additional harness tests, and strict all-target Clippy. + +## Reproduce + +Run `make candidates` from `promql-compliance/runner`. The committed Compose +configuration enables durable state required by maintenance-stage candidates. +The initial in-memory run failed three candidate trials with +`immutable maintenance requires durable input state`. The isolated replay and +complete durable rerun pass; no candidate was skipped or tolerance relaxed. + +## Provenance and artifacts + +- Runner code built at `b2d521fe8ca6eed95d1f3e8431a7b1df670b7658`. +- Durable Compose fixture committed at `d393b1bf`; local run mounted the existing + binary instead of rebuilding the image. The mounted data-plane binary was + built at `ed944356`; its runtime source is unchanged in the tested revision. +- Planner: `176c1bd565e0c400f9a35996c2bd65de32475e72`. +- Data-plane SHA256: `32a89c8f2091633a629e7fd1a49dc264447069bb95cfd123554f1db9ad8311e8`. +- Runner SHA256: `3bb04c99d9e382d54945efbefb316edce9305b7ea8c2bed1643da0ee80e6d5f0`. + +Generated snapshots, selected plans, installations, responses and service logs +are kept outside version control. Reproduce them with the command above. diff --git a/docs/evaluation/dataset-bound-sds-2026-09-28/README.md b/docs/evaluation/dataset-bound-sds-2026-09-28/README.md new file mode 100644 index 000000000..c6daa7f61 --- /dev/null +++ b/docs/evaluation/dataset-bound-sds-2026-09-28/README.md @@ -0,0 +1,67 @@ +# Dataset-bound SDS acceptance evidence + +The implementation follows the logical dataset, same-version recovery and query-wide +consistency contracts in #737. This run uses synthetic candidate prices; online ERP, +cross-version state adoption and ad-hoc SDS discovery remain deferred. + +## Accepted behavior + +| Contract | Evidence | +| --- | --- | +| Same expression, different dataset | Different semantic definition IDs; mismatched or absent installed input identity is rejected. | +| Same dataset, relocated input | Changing the collector binding preserves the semantic definitions. | +| Same-version recovery | A real process restarts from disk and returns 15 without re-ingestion. | +| New-version warm-up | Identical definitions do not grant access to old state; fresh input warms the new version and returns 30. | +| Query-wide consistency | Publishing between two query branches invalidates the whole result. | +| Cost isolation | Synthetic ranking rejects quotes from another dataset. | + +Planner carries an explicit namespace/dataset in semantic fragment version 2. +Backend snapshots use version 3 and installed catalogs use schema 6. One input +channel has one declared logical dataset; the trusted source authority supplies +its identity. The implementation does not infer tenants from metric names or bytes. + +## Data-plane sweep + +Every admitted fixture candidate was made cheapest using synthetic prices, +compiled into deployment plans, installed, ingested through Remote Write, drained, +and compared with Prometheus. Ensembles share one installation per candidate. + +| Workload | Candidate runs | +| --- | --- | +| query:spatial-sum | 1 | +| query:spatial-topk | 1 | +| query:spatial-quantile | 1 | +| query:temporal-sum | 1 | +| query:temporal-quantile | 1 | +| query:temporal-rate | 1 | +| query:grouped-rate | 3 | +| query:grouped-temporal-sum | 2 | +| query:topk-rate | 2 | +| query:quantile-ratio | 1 | +| shared-rate | 2 | +| shared-quantiles | 1 | +| full-ensemble | 3 | + +Total: 20 candidate runs and 52 query comparisons; all passed. + +Admission reports retain rejected candidates. This covers the admitted fixture +inventory, not every theoretical plan, production cost optimality, or human review. + +## Validation and provenance + +- Backend runtime and runner built at `ea809e00c88729a7a623d62020ba94166b5045fe`. +- Later commits before this report only propagate documentation and plan exports. +- Planner code pin: `bccc837f5caf21c888b2cce73c64a1be283b679a`. +- Three Level 1 tests and the full synthetic Level 2 ranking test passed. +- Focused binding, snapshot-version, adoption, recovery, query-fence and installation tests passed. +- Strict all-target Clippy passed for control_plane, data_plane and promql-compliance. +- Deferred test stack passed an all-target compile check; its production-cost acceptance was not rerun. +- Cross-version metadata regression failed before the fix; generated logs are kept outside version control. +- Existing runtime generation filtering was already strict; the fix rejects the obsolete metadata marker. + +- `data_plane` SHA256: `b3d4d37f6520670955d2b3f9fd2afa208ccb001ec108370a542f47727e5d5206`. +- `differential-runner` SHA256: `0b8a87d439b6943e178672f01a544adf90ea2644f3065dc2bb9bad1e7941b6e2`. + +Run `make candidates` in `promql-compliance/runner` to reproduce the sweep. +The run used built binaries and durable state. Generated logs, Compose captures +and detailed result artifacts are kept outside version control. diff --git a/docs/evaluation/execution-2026-09-28/README.md b/docs/evaluation/execution-2026-09-28/README.md new file mode 100644 index 000000000..c07a1298f --- /dev/null +++ b/docs/evaluation/execution-2026-09-28/README.md @@ -0,0 +1,21 @@ +# Installed-plan execution: 2026-09-28 + +All seven container runs passed against Prometheus: single-rate-temporal, +sparse-checkout-temporal, aggregations, aggregations-dense-cadence, issue-702, +issue-702-one-second and issue-754. See the retained summary and +[provenance](provenance.json) for the exact runner/runtime +revisions. The execution harness has been moved without changing its behavior. + +The runner compiled a plan offline, validated its local execution path, and +installed that exact typed plan. Server startup did not re-plan. Reports retain +backend/reference results and execution provenance. Inputs are finite generated +fixtures; their accuracy contracts and analytical costs are not production ERP +measurements. This proves execution behavior for those plans, not that cost +selection matches production reality. + +Reproduce with `promql-compliance/runner/Makefile` and the dataset/suite pairs +listed above; see [runner usage](../../../promql-compliance/README.md). The recorded +run used the debug binary mounted in the runtime container. No human review, +performance benefit, or Level 3 selection approval is implied. + +Generated detailed reports are kept outside version control. diff --git a/docs/evaluation/execution-2026-09-28/provenance.json b/docs/evaluation/execution-2026-09-28/provenance.json new file mode 100644 index 000000000..1f7ff9a2e --- /dev/null +++ b/docs/evaluation/execution-2026-09-28/provenance.json @@ -0,0 +1,12 @@ +{ + "planner_revision": "176c1bd565e0c400f9a35996c2bd65de32475e72", + "runner_revision": "13e90c8d", + "data_plane_revision": "ed944356", + "scope": "finite fixture differential execution, not production selection or real ERP calibration", + "reports": [ + { + "file": "summary.json", + "sha256": "3f7bad5dda4fdb08b235afaca9adacaf825262598fa8d3cd752437b1565bcac3" + } + ] +} diff --git a/docs/evaluation/execution-2026-09-28/summary.json b/docs/evaluation/execution-2026-09-28/summary.json new file mode 100644 index 000000000..cd9e4543d --- /dev/null +++ b/docs/evaluation/execution-2026-09-28/summary.json @@ -0,0 +1,61 @@ +{ + "cases": [ + { + "asapQuery": 330, + "dataset": "aggregations", + "passed": true, + "prometheusFallback": 0, + "queries": 66, + "suite": "aggregations" + }, + { + "asapQuery": 45, + "dataset": "aggregations-dense-cadence", + "passed": true, + "prometheusFallback": 0, + "queries": 9, + "suite": "issue-702" + }, + { + "asapQuery": 330, + "dataset": "aggregations-dense-cadence", + "passed": true, + "prometheusFallback": 0, + "queries": 66, + "suite": "aggregations" + }, + { + "asapQuery": 16, + "dataset": "issue-702-one-second", + "passed": true, + "prometheusFallback": 0, + "queries": 8, + "suite": "issue-702-one-second" + }, + { + "asapQuery": 30, + "dataset": "issue-754", + "passed": true, + "prometheusFallback": 0, + "queries": 10, + "suite": "issue-754" + }, + { + "asapQuery": 4, + "dataset": "single-rate", + "passed": true, + "prometheusFallback": 0, + "queries": 1, + "suite": "temporal" + }, + { + "asapQuery": 4, + "dataset": "sparse-checkout", + "passed": true, + "prometheusFallback": 0, + "queries": 1, + "suite": "temporal" + } + ], + "passed": true +} \ No newline at end of file diff --git a/promql-compliance/README.md b/promql-compliance/README.md new file mode 100644 index 000000000..829da9ff4 --- /dev/null +++ b/promql-compliance/README.md @@ -0,0 +1,149 @@ +> Current milestone: #728 candidate structure → #742 synthetic ranking → #775 +> selected-plan execution. Online ERP collection and runtime replanning are deferred. +> This suite establishes execution correctness, not measured cost optimality. + +# PromQL compliance suite + +The Rust `promql-compliance` workspace crate 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 BuildKit Docker build can take several minutes; subsequent builds reuse Cargo caches. 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 runner adds `servedBy` from execution provenance headers: + +- `asap_query`: both `X-ASAP-Execution: warm` and `X-ASAP-Execution-Detail: asap` confirm local execution. +- `hybrid`: includes external exact work and fails strict local acceptance. +- `prometheus_fallback`: fallback or missing/invalid local evidence; it cannot pass strict acceptance. + +## 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 Physical DAG. The runner derives a deployment planning input from the fixture and suite; Planner exposes physical candidates; Backend selects using explicit synthetic quotes and local execution feasibility. + +The suite query expression is the workload query. Dataset metric names must cover +its named sources. `runner/src/planning.rs` builds a typed backend planning input: +a recurrence and phase covering the actual evaluation timestamps, with explicit +1% epsilon/delta accuracy. The replay declares the historical retention needed +between its latest input and earliest evaluation. Complete finite fixture bounds +supply scoped operand-domain proofs; these are not inferred guarantees about live data. Differential fixtures derive total arrival rate from each series' replay span +and use the largest per-series sample gap as the declared cadence bound. +Benefit fixtures require uniform source cadence. Both paths retain the actual +series count and sample volume; no synthetic one-second cadence is supplied. + +The differential runner invokes the Rust control-plane compiler directly. It quotes +every component at a deterministic unit price of 1.0 under +`synthetic-execution-fixture-v1`; non-local candidates are infeasible for this +fixture. These fake prices select a plan through the production quote path, without +ERP measurements or runtime feedback. The runner validates the synthetic cost +report and local execution before installing that exact deployment plan. Backend +startup does not repeat candidate selection. Planning failures produce a JSON +report. The input snapshot, selected plan, and typed `*.install.json` publication +are saved beside successful planning reports. Docker Compose is needed only for live +services; fixture, protocol, cost and comparison tests run with Cargo: + +```sh +cargo test --locked -p promql-compliance +cargo clippy --locked -p promql-compliance --all-targets -- -D warnings +``` + +There is no Go toolchain, runner, module or generated Go code in this harness. + +### Candidate execution matrix + +Run `make candidates` in `promql-compliance/runner` to test each issue-754 query, +the shared-rate and shared-quantile ensembles, and the complete ten-query workload. +For every admitted executable candidate, the harness changes only synthetic quotes +to make that candidate cheapest, asserts that production selection chose it, +compiles and installs its deployment plan, ingests the fixture once, and compares +every query in that workload with Prometheus. A failed candidate fails the matrix. + +The report retains candidate admissions, quoted snapshots, selected plans and typed +installations. Coverage means all candidates exposed and admitted for this fixture; +it does not claim an exhaustive Cartesian product of all possible DAGs. Candidates +rejected by accuracy admission or deployment binding remain recorded as rejected. +Online ERP measurements and runtime replanning are outside this test. + +The Compose deployment enables Remote Write and durable summary storage. Some +candidates read persisted intermediate state during precompute; an in-memory-only +profile cannot execute them. State stays within the fresh container for a trial, +and teardown discards it before installing the next candidate. diff --git a/promql-compliance/datasets/aggregations-dense-cadence.yaml b/promql-compliance/datasets/aggregations-dense-cadence.yaml new file mode 100644 index 000000000..f7f719f90 --- /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 000000000..31c987978 --- /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/issue-702-one-second.yaml b/promql-compliance/datasets/issue-702-one-second.yaml new file mode 100644 index 000000000..a7228c9f6 --- /dev/null +++ b/promql-compliance/datasets/issue-702-one-second.yaml @@ -0,0 +1,23 @@ +name: issue-702-one-second +# Mirrors data_plane/tests/support/issue_701_702_process.rs: three series, +# one sample per second from 1 through 961, factor * (1 + second % 31). +series: + - metric: issue701_data + labels: {pod: a, job: api} + generated_samples: &factor_one + start_offset_seconds: 1 + end_offset_seconds: 961 + step_seconds: 1 + multiplier: 1 + base: 1 + modulo: 31 + - metric: issue701_data + labels: {pod: b, job: api} + generated_samples: + <<: *factor_one + multiplier: 2 + - metric: issue701_data + labels: {pod: c, job: db} + generated_samples: + <<: *factor_one + multiplier: 3 diff --git a/promql-compliance/datasets/single-rate.yaml b/promql-compliance/datasets/single-rate.yaml new file mode 100644 index 000000000..8bc52bf9e --- /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 000000000..d9b71da87 --- /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 000000000..8426afbcd --- /dev/null +++ b/promql-compliance/docker-compose.yml @@ -0,0 +1,23 @@ +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"] + healthcheck: + test: ["CMD", "/bin/wget", "-q", "-O", "/dev/null", "http://localhost:9090/-/healthy"] + interval: 1s + timeout: 2s + retries: 60 + data-plane: + depends_on: + prometheus: + condition: service_healthy + build: + context: .. + dockerfile: data_plane/Dockerfile + command: ["--enable-remote-write", "--physical-plan", "/config/physical-plan.json", "--prometheus-server", "http://prometheus:9090", "--http-port", "9091", "--output-dir", "/tmp/asap", "--persistence-enabled", "--persistence-dir", "/tmp/asap/state", "--persistence-delete-older-than-secs", "0"] + volumes: + - ${ASAP_PHYSICAL_PLAN}:/config/physical-plan.json:ro + ports: ["${BACKEND_PORT:-19091}:9091"] diff --git a/promql-compliance/runner/Cargo.toml b/promql-compliance/runner/Cargo.toml new file mode 100644 index 000000000..5168ba025 --- /dev/null +++ b/promql-compliance/runner/Cargo.toml @@ -0,0 +1,23 @@ +[package] +name = "promql-compliance" +version.workspace = true +edition.workspace = true +publish = false + +[dependencies] +anyhow.workspace = true +clap = { workspace = true, features = ["derive"] } +serde.workspace = true +serde_json.workspace = true +serde_yaml.workspace = true +tokio.workspace = true +reqwest.workspace = true +control_plane = { path = "../../control_plane" } +asap_types.workspace = true +planner-types.workspace = true +prost = "0.13" +snap = "1.1" +regex = "1" + +[dev-dependencies] +tempfile = "3" diff --git a/promql-compliance/runner/Makefile b/promql-compliance/runner/Makefile new file mode 100644 index 000000000..4d523cb6a --- /dev/null +++ b/promql-compliance/runner/Makefile @@ -0,0 +1,36 @@ +.PHONY: test run run-all candidates + +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 \ + issue-702:../datasets/aggregations-dense-cadence.yaml:../suites/issue-702.yaml \ + issue-702-one-second:../datasets/issue-702-one-second.yaml:../suites/issue-702-one-second.yaml \ + issue-754:../datasets/issue-754.yaml:../suites/issue-754.yaml + +test: + cargo test --locked -p promql-compliance + +run: + cargo run --locked -p promql-compliance --bin 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 ! cargo run --locked -p promql-compliance --bin 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; cargo run --locked -p promql-compliance --bin report-card -- --reports-dir "$(REPORT_DIR)"; exit $$result + +# Select, install and execute every admitted candidate for singles and ensembles. +candidates: + cargo run --locked -p promql-compliance --bin differential-runner -- --all-candidates --dataset ../datasets/issue-754.yaml --suite ../suites/issue-754.yaml --reference-url $(REFERENCE_URL) --test-url $(BACKEND_URL) --compose-file $(COMPOSE_FILE) --logs-dir "$(LOGS_DIR)/candidates" --output "$(REPORT_DIR)/candidates.json" diff --git a/promql-compliance/runner/sql/counter.sql b/promql-compliance/runner/sql/counter.sql new file mode 100644 index 000000000..07ff9d37c --- /dev/null +++ b/promql-compliance/runner/sql/counter.sql @@ -0,0 +1,37 @@ +, +ordered AS ( + SELECT series_id, label_0, ts_ms, value, + row_number() OVER (PARTITION BY series_id ORDER BY ts_ms) AS sample_index, + lag(value, 1, 0.) OVER (PARTITION BY series_id ORDER BY ts_ms) AS previous_value + FROM window_samples +), corrected AS ( + SELECT series_id, label_0, count() AS sample_count, + min(ts_ms) AS first_ms, max(ts_ms) AS last_ms, + argMin(value, ts_ms) AS first_value, argMax(value, ts_ms) AS last_value, + sum(if(sample_index > 1 AND value < previous_value, previous_value, 0.)) AS reset_correction + FROM ordered GROUP BY series_id, label_0 HAVING sample_count >= 2 +), durations AS ( + SELECT *, last_value - first_value + reset_correction AS corrected_delta, + (last_ms - first_ms) / 1000. AS sampled_seconds, + (first_ms - (t_ms - window_ms)) / 1000. AS gap_start_seconds, + (t_ms - last_ms) / 1000. AS gap_end_seconds + FROM corrected +), thresholds AS ( + SELECT *, sampled_seconds / (sample_count - 1) AS average_gap_seconds FROM durations +), adjusted AS ( + SELECT *, + if(gap_start_seconds >= 1.1 * average_gap_seconds, average_gap_seconds / 2, gap_start_seconds) AS start_seconds, + if(gap_end_seconds >= 1.1 * average_gap_seconds, average_gap_seconds / 2, gap_end_seconds) AS end_seconds + FROM thresholds +), extrapolated AS ( + SELECT *, if(corrected_delta > 0 AND first_value >= 0, + least(start_seconds, sampled_seconds * first_value / corrected_delta), + start_seconds) AS zero_adjusted_start_seconds FROM adjusted +), per_series_counter AS ( + SELECT series_id, label_0, + corrected_delta * (sampled_seconds + zero_adjusted_start_seconds + end_seconds) + / sampled_seconds AS increase_value, + corrected_delta * (sampled_seconds + zero_adjusted_start_seconds + end_seconds) + / sampled_seconds / (window_ms / 1000.) AS rate_value + FROM extrapolated +) \ No newline at end of file diff --git a/promql-compliance/runner/sql/prefix.sql b/promql-compliance/runner/sql/prefix.sql new file mode 100644 index 000000000..7a3abb780 --- /dev/null +++ b/promql-compliance/runner/sql/prefix.sql @@ -0,0 +1,9 @@ +WITH {evaluation_ms} AS t_ms, {window_ms} AS window_ms, +instant_samples AS ( + SELECT series_id, label_0, argMax(value, ts_ms) AS value + FROM samples WHERE ts_ms <= t_ms AND ts_ms >= t_ms - 300000 + GROUP BY series_id, label_0 +), window_samples AS ( + SELECT series_id, label_0, ts_ms, value FROM samples + WHERE ts_ms > t_ms - window_ms AND ts_ms <= t_ms +) \ No newline at end of file diff --git a/promql-compliance/runner/src/bin/differential-runner.rs b/promql-compliance/runner/src/bin/differential-runner.rs new file mode 100644 index 000000000..62281110b --- /dev/null +++ b/promql-compliance/runner/src/bin/differential-runner.rs @@ -0,0 +1,5 @@ +use clap::Parser; +#[tokio::main] +async fn main() -> anyhow::Result<()> { + promql_compliance::runner::run(promql_compliance::runner::Args::parse(), false).await +} diff --git a/promql-compliance/runner/src/bin/report-card.rs b/promql-compliance/runner/src/bin/report-card.rs new file mode 100644 index 000000000..9580975de --- /dev/null +++ b/promql-compliance/runner/src/bin/report-card.rs @@ -0,0 +1,9 @@ +use clap::Parser; +#[derive(Parser)] +struct Args { + #[arg(long)] + reports_dir: std::path::PathBuf, +} +fn main() -> anyhow::Result<()> { + promql_compliance::runner::report_card(&Args::parse().reports_dir) +} diff --git a/promql-compliance/runner/src/compare.rs b/promql-compliance/runner/src/compare.rs new file mode 100644 index 000000000..955bb759c --- /dev/null +++ b/promql-compliance/runner/src/compare.rs @@ -0,0 +1,150 @@ +use crate::input::Policy; +use anyhow::{bail, ensure, Context, Result}; +use serde_json::{json, Value}; + +pub fn equal_number(a: f64, b: f64, policy: &Policy) -> bool { + if a.is_nan() || b.is_nan() { + return a.is_nan() && b.is_nan(); + } + if a.is_infinite() || b.is_infinite() { + return a == b; + } + let relative = policy + .value_tolerance + .as_ref() + .and_then(|t| t.relative) + .unwrap_or(0.); + let absolute = policy + .value_tolerance + .as_ref() + .and_then(|t| t.absolute) + .unwrap_or(0.); + (a - b).abs() <= absolute + relative * a.abs().max(b.abs()) +} +#[derive(Debug, PartialEq)] +pub struct Point { + pub labels: String, + pub timestamp: i64, + pub value: f64, +} +pub fn normalize(response: &Value) -> Result> { + ensure!(response["status"] == "success", "query failed: {response}"); + let data = &response["data"]; + let mut points = Vec::new(); + match data["resultType"].as_str() { + Some("vector" | "matrix") => { + for s in data["result"].as_array().context("invalid result series")? { + let labels: std::collections::BTreeMap = + serde_json::from_value(s["metric"].clone())?; + let labels = serde_json::to_string(&labels)?; + let values = if data["resultType"] == "vector" { + vec![&s["value"]] + } else { + s["values"] + .as_array() + .context("invalid matrix")? + .iter() + .collect() + }; + for pair in values { + let (timestamp, value) = point(pair)?; + points.push(Point { + labels: labels.clone(), + timestamp, + value, + }); + } + } + } + Some("scalar") => { + let (timestamp, value) = point(&data["result"])?; + points.push(Point { + labels: String::new(), + timestamp, + value, + }); + } + _ => bail!("unsupported numeric result type"), + } + points.sort_by(|a, b| a.labels.cmp(&b.labels).then(a.timestamp.cmp(&b.timestamp))); + ensure!( + !points + .windows(2) + .any(|p| p[0].labels == p[1].labels && p[0].timestamp == p[1].timestamp), + "duplicate response sample" + ); + Ok(points) +} +fn point(value: &Value) -> Result<(i64, f64)> { + let pair = value.as_array().context("invalid sample")?; + ensure!(pair.len() == 2, "sample needs timestamp and value"); + Ok(( + crate::input::offset_ms(pair[0].as_f64().context("invalid timestamp")?)?, + pair[1].as_str().context("invalid sample value")?.parse()?, + )) +} +pub fn compare(left: &Value, right: &Value, policy: &Policy) -> Result<()> { + ensure!( + left["status"] == "success" && right["status"] == "success", + "failed API response" + ); + ensure!( + left["data"]["resultType"] == right["data"]["resultType"], + "result type mismatch" + ); + if left["data"]["resultType"] == "string" { + ensure!( + left["data"]["result"] == right["data"]["result"], + "string/timestamp mismatch" + ); + return Ok(()); + } + let a = normalize(left)?; + let b = normalize(right)?; + ensure!( + a.len() == b.len(), + "sample count mismatch {} != {}", + a.len(), + b.len() + ); + for (a, b) in a.iter().zip(b.iter()) { + ensure!( + a.labels == b.labels && a.timestamp == b.timestamp, + "labels/timestamp mismatch: {a:?} != {b:?}" + ); + ensure!( + equal_number(a.value, b.value, policy), + "value mismatch: {} != {}", + a.value, + b.value + ); + } + Ok(()) +} +pub fn parity(range: &Value, instant: &Value, at: i64, policy: &Policy) -> Result<()> { + ensure!( + range["data"]["resultType"] == "matrix" && instant["data"]["resultType"] == "vector", + "invalid range/instant types" + ); + let a = normalize(range)? + .into_iter() + .filter(|p| p.timestamp == at) + .collect::>(); + let b = normalize(instant)?; + ensure!(a.len() == b.len(), "range/instant sample count mismatch"); + for (a, b) in a.iter().zip(b.iter()) { + ensure!( + a.labels == b.labels + && a.timestamp == b.timestamp + && equal_number(a.value, b.value, policy), + "range/instant mismatch" + ); + } + Ok(()) +} +pub fn outcome(result: Result<()>) -> Value { + match result { + Ok(()) => json!({"passed":true}), + Err(e) => json!({"passed":false,"diff":format!("{e:#}")}), + } +} diff --git a/promql-compliance/runner/src/compose.rs b/promql-compliance/runner/src/compose.rs new file mode 100644 index 000000000..19e34114c --- /dev/null +++ b/promql-compliance/runner/src/compose.rs @@ -0,0 +1,122 @@ +use anyhow::{bail, ensure, Result}; +use std::{path::PathBuf, process::Command}; + +pub struct Compose { + pub files: Vec, + pub project: String, + pub installation: PathBuf, + pub logs: PathBuf, + pub keep: bool, + started: bool, +} +impl Compose { + pub fn new( + files: Vec, + project: String, + installation: PathBuf, + logs: PathBuf, + keep: bool, + ) -> Self { + Self { + files, + project, + installation, + logs, + keep, + started: false, + } + } + fn command(&self) -> Command { + let mut c = Command::new("docker"); + c.arg("compose").args(["--project-name", &self.project]); + for file in &self.files { + c.arg("--file").arg(file); + } + c.env("ASAP_PHYSICAL_PLAN", &self.installation); + c + } + fn run(&self, args: &[&str]) -> Result { + let output = self.command().args(args).output()?; + std::fs::create_dir_all(&self.logs)?; + use std::io::Write; + let mut log = std::fs::OpenOptions::new() + .append(true) + .create(true) + .open(self.logs.join("lifecycle.log"))?; + log.write_all(&output.stdout)?; + log.write_all(&output.stderr)?; + ensure!( + output.status.success(), + "Compose {} failed: {}", + args.join(" "), + String::from_utf8_lossy(&output.stderr) + ); + Ok(String::from_utf8(output.stdout)?) + } + pub fn start(&mut self, benefit: bool) -> Result<()> { + if self.files.is_empty() { + return Ok(()); + } + self.run(&["down", "--volumes", "--remove-orphans"])?; + self.started = true; + self.run(&["up", "-d", "--build", "prometheus", "data-plane"])?; + if benefit { + self.run(&["up", "-d", "clickhouse", "victoria"])?; + } + Ok(()) + } + pub fn usage(&self, service: &str) -> Result<(u64, u64)> { + ensure!(!self.files.is_empty(), "resource measurement needs Compose"); + let id = self.run(&["ps", "-q", service])?; + ensure!( + !id.trim().is_empty() && id.trim().lines().count() == 1, + "expected one {service} container" + ); + let read = |path: &str| -> Result { + let out = Command::new("docker") + .args(["exec", id.trim(), "cat", path]) + .output()?; + ensure!( + out.status.success(), + "read cgroup {path}: {}", + String::from_utf8_lossy(&out.stderr) + ); + Ok(String::from_utf8(out.stdout)?) + }; + let cpu = read("/sys/fs/cgroup/cpu.stat")?; + let usage = parse_cpu_stat(&cpu)?; + let peak = read("/sys/fs/cgroup/memory.peak")?.trim().parse()?; + Ok((usage, peak)) + } + pub fn finish(&mut self) -> Result<()> { + if !self.started { + return Ok(()); + } + let logs = self.run(&["logs", "--no-color"]); + if let Ok(logs) = logs { + std::fs::write(self.logs.join("compose.log"), logs)?; + } + if self.keep { + self.started = false; + return Ok(()); + } + self.run(&["down", "--volumes", "--remove-orphans"])?; + self.started = false; + Ok(()) + } +} +impl Drop for Compose { + fn drop(&mut self) { + if let Err(e) = self.finish() { + eprintln!("Compose cleanup failed: {e:#}"); + } + } +} +pub fn parse_cpu_stat(text: &str) -> Result { + for line in text.lines() { + if let Some(value) = line.strip_prefix("usage_usec ") { + return Ok(value.parse()?); + } + } + bail!("missing cgroup-v2 usage_usec") +} diff --git a/promql-compliance/runner/src/input.rs b/promql-compliance/runner/src/input.rs new file mode 100644 index 000000000..283dff56a --- /dev/null +++ b/promql-compliance/runner/src/input.rs @@ -0,0 +1,285 @@ +use anyhow::{ensure, Context, Result}; +use serde::{Deserialize, Serialize}; +use std::{ + collections::{BTreeMap, BTreeSet}, + path::Path, +}; + +#[derive(Clone, Debug, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct Suite { + pub name: String, + #[serde(default)] + pub comparison_defaults: Policy, + pub queries: Vec, +} +#[derive(Clone, Debug, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct Query { + pub name: String, + pub expr: String, + #[serde(default)] + pub instant_offsets_seconds: Vec, + pub range: Option, + pub comparison: Option, +} +#[derive(Clone, Debug, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct Range { + pub start_offset_seconds: f64, + pub end_offset_seconds: f64, + pub step_seconds: f64, +} +#[derive(Clone, Debug, Default, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct Policy { + pub value_tolerance: Option, +} +#[derive(Clone, Debug, Default, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct Tolerance { + pub relative: Option, + pub absolute: Option, +} +impl Query { + pub fn policy(&self, defaults: &Policy) -> Policy { + let base = defaults.value_tolerance.clone().unwrap_or_default(); + let own = self + .comparison + .as_ref() + .and_then(|p| p.value_tolerance.clone()) + .unwrap_or_default(); + Policy { + value_tolerance: Some(Tolerance { + relative: own.relative.or(base.relative), + absolute: own.absolute.or(base.absolute), + }), + } + } +} +impl Suite { + pub fn load(path: &Path) -> Result { + Self::parse(&std::fs::read_to_string(path)?) + } + pub fn parse(text: &str) -> Result { + let suite: Self = yaml(text)?; + ensure!( + !suite.name.is_empty() && !suite.queries.is_empty(), + "empty suite" + ); + let mut names = BTreeSet::new(); + for q in &suite.queries { + ensure!( + !q.name.is_empty() && !q.expr.is_empty() && names.insert(&q.name), + "empty/duplicate query" + ); + ensure!( + !q.instant_offsets_seconds.is_empty() || q.range.is_some(), + "query {} has no evaluations", + q.name + ); + if let Some(r) = &q.range { + ensure!( + [r.start_offset_seconds, r.end_offset_seconds, r.step_seconds] + .iter() + .all(|n| n.is_finite()) + && r.end_offset_seconds > r.start_offset_seconds + && r.step_seconds > 0., + "invalid query range" + ); + } + for t in &q.instant_offsets_seconds { + ensure!( + t.is_finite() + && q.range.as_ref().is_none_or( + |r| *t >= r.start_offset_seconds && *t <= r.end_offset_seconds + ), + "invalid instant time" + ); + } + let t = q + .policy(&suite.comparison_defaults) + .value_tolerance + .unwrap(); + ensure!( + [t.relative, t.absolute] + .into_iter() + .flatten() + .all(|n| n.is_finite() && n >= 0.), + "invalid tolerance" + ); + } + Ok(suite) + } +} +#[derive(Clone, Debug, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct Dataset { + pub name: String, + pub series: Vec, +} +#[derive(Clone, Debug, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct Series { + pub metric: String, + #[serde(default)] + pub labels: BTreeMap, + #[serde(default)] + pub samples: Vec, + pub generated_samples: Option, +} +#[derive(Clone, Debug, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct Sample { + pub offset_seconds: f64, + pub value: f64, +} +#[derive(Clone, Debug, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct Generated { + pub start_offset_seconds: f64, + pub end_offset_seconds: f64, + pub step_seconds: f64, + pub multiplier: f64, + pub base: f64, + pub modulo: f64, +} +impl Dataset { + pub fn load(path: &Path) -> Result { + Self::parse(&std::fs::read_to_string(path)?) + } + pub fn parse(text: &str) -> Result { + let mut data: Self = yaml(text)?; + ensure!( + !data.name.is_empty() && !data.series.is_empty(), + "empty dataset" + ); + let mut seen = BTreeSet::new(); + for s in &mut data.series { + ensure!( + !s.metric.is_empty() + && !s.labels.contains_key("__name__") + && seen.insert((s.metric.clone(), s.labels.clone())), + "invalid/duplicate series" + ); + if let Some(g) = s.generated_samples.take() { + ensure!( + s.samples.is_empty(), + "explicit and generated samples both supplied" + ); + ensure!( + [ + g.start_offset_seconds, + g.end_offset_seconds, + g.step_seconds, + g.multiplier, + g.base, + g.modulo + ] + .iter() + .all(|n| n.is_finite()) + && g.end_offset_seconds >= g.start_offset_seconds + && g.step_seconds > 0. + && g.modulo > 0., + "invalid generated samples" + ); + let count = (g.end_offset_seconds - g.start_offset_seconds) / g.step_seconds; + ensure!( + count < 10_000_000. && (count - count.round()).abs() < 1e-9, + "invalid generated sample grid" + ); + s.samples = (0..=count.round() as usize) + .map(|i| { + let offset = g.start_offset_seconds + i as f64 * g.step_seconds; + Sample { + offset_seconds: offset, + value: g.multiplier * (g.base + offset % g.modulo), + } + }) + .collect(); + } + ensure!(!s.samples.is_empty(), "empty series"); + let mut previous = None; + for sample in &s.samples { + ensure!( + sample.offset_seconds.is_finite() && sample.value.is_finite(), + "nonfinite sample" + ); + let ms = offset_ms(sample.offset_seconds)?; + ensure!( + previous.is_none_or(|p| ms > p), + "samples collide or are unordered at millisecond precision" + ); + previous = Some(ms); + } + } + Ok(data) + } + /// Complete replay bounds: total source rate and largest per-series gap. + pub fn replay_demand(&self) -> Result<(f64, u64, usize)> { + let mut rate = 0.; + let mut cadence = 0; + let mut count = 0; + for series in &self.series { + ensure!( + series.samples.len() >= 2, + "replay cadence needs two samples per series" + ); + let first = offset_ms(series.samples.first().unwrap().offset_seconds)?; + let last = offset_ms(series.samples.last().unwrap().offset_seconds)?; + rate += (series.samples.len() - 1) as f64 * 1000. / (last - first) as f64; + for pair in series.samples.windows(2) { + let gap = offset_ms(pair[1].offset_seconds)? - offset_ms(pair[0].offset_seconds)?; + ensure!(gap > 0, "replay timestamps must increase"); + cadence = cadence.max(gap as u64); + } + count += series.samples.len(); + } + ensure!( + cadence > 0 && rate.is_finite() && rate > 0., + "empty replay demand" + ); + Ok((rate, cadence, count)) + } + + pub fn uniform_demand(&self) -> Result<(f64, u64, usize)> { + let mut cadence = None; + let mut count = 0; + for s in &self.series { + ensure!( + s.samples.len() >= 2, + "source {} needs two samples for cadence", + s.metric + ); + for pair in s.samples.windows(2) { + let step = offset_ms(pair[1].offset_seconds)? - offset_ms(pair[0].offset_seconds)?; + ensure!( + step > 0 && cadence.is_none_or(|p| p == step), + "benefit fixture needs uniform source cadence" + ); + cadence = Some(step); + } + count += s.samples.len(); + } + let ms = cadence.context("empty data population")? as u64; + Ok((self.series.len() as f64 * 1000. / ms as f64, ms, count)) + } +} +pub fn offset_ms(seconds: f64) -> Result { + let ms = (seconds * 1000.).round(); + ensure!( + ms.is_finite() && ms > i64::MIN as f64 && ms < i64::MAX as f64, + "timestamp overflow" + ); + Ok(ms as i64) +} +pub fn at_ms(base: i64, seconds: f64) -> Result { + base.checked_add(offset_ms(seconds)?) + .context("timestamp overflow") +} + +fn yaml(text: &str) -> Result { + let mut value: serde_yaml::Value = serde_yaml::from_str(text)?; + value.apply_merge()?; + Ok(serde_yaml::from_value(value)?) +} diff --git a/promql-compliance/runner/src/lib.rs b/promql-compliance/runner/src/lib.rs new file mode 100644 index 000000000..eebbf5f44 --- /dev/null +++ b/promql-compliance/runner/src/lib.rs @@ -0,0 +1,7 @@ +pub mod compare; +pub mod compose; +pub mod input; +pub mod planning; +pub mod runner; +pub mod sql; +pub mod transport; diff --git a/promql-compliance/runner/src/planning.rs b/promql-compliance/runner/src/planning.rs new file mode 100644 index 000000000..78f4ae1a2 --- /dev/null +++ b/promql-compliance/runner/src/planning.rs @@ -0,0 +1,466 @@ +use crate::input::{at_ms, Dataset, Suite}; +use anyhow::{ensure, Context, Result}; +use asap_types::query_plan::{residual::ResidualQueryOperator, QueryPlanNode}; +use control_plane::physical::{ + compiler::{BackendLocalPlanningInput, CompiledPhysicalPlan}, + workload_cost::CandidateEvaluationStatus, +}; +use serde_json::{json, Value}; + +pub fn snapshot( + suite: &Suite, + dataset: &Dataset, + now: u64, + base_time_ms: i64, + benefit: bool, +) -> Result { + let mut value: Value = serde_json::from_str(include_str!( + "../../../docs/examples/asapquery-planning-snapshot.json" + ))?; + let queries: Vec<_> = suite.queries.iter().map(|q| { + let ms=q.range.as_ref().map(|r| r.step_seconds*1000.).unwrap_or(60_000.); + ensure!(ms.fract()==0. && ms>0. && ms<=u32::MAX as f64,"query recurrence needs positive integral milliseconds"); + let mut interval = ms as u64; + let times = q.instant_offsets_seconds.iter().copied() + .chain(q.range.iter().map(|r| r.start_offset_seconds)) + .map(|offset| at_ms(base_time_ms, offset)).collect::>>()?; + let first = *times.first().context("query has no evaluation times")?; + for time in times { + interval = gcd(interval, time.abs_diff(first)); + } + let phase = first.rem_euclid(interval as i64) as u64; + Ok(json!({"query":q.expr,"demand":{"fixed_interval_at":{"interval":interval,"evaluation_phase":phase}}, + "requirements":{"accuracy":{"explicit":{"EpsilonDelta":{"epsilon":0.01,"delta":0.01}}},"response_latency":"unspecified"}, + "predictability":{"predictable":{"known_at":null}},"time_selection":{"scope":"real_time","lookback":null,"as_of":null}})) + }).collect::>()?; + value["query_workload"]["repeating_queries"] = json!(queries); + value["implementation"]["require_backend_local_execution"] = json!(true); + let newest = dataset + .series + .iter() + .flat_map(|s| &s.samples) + .map(|s| at_ms(base_time_ms, s.offset_seconds)) + .collect::>>()? + .into_iter() + .max() + .context("empty replay")?; + let earliest = suite + .queries + .iter() + .flat_map(|q| { + q.instant_offsets_seconds + .iter() + .copied() + .chain(q.range.iter().map(|r| r.start_offset_seconds)) + }) + .map(|t| at_ms(base_time_ms, t)) + .collect::>>()? + .into_iter() + .min() + .context("no evaluation times")?; + value["implementation"]["query_staleness_margin_ms"] = + json!(newest.saturating_sub(earliest).max(0) as u64); + + value["implementation"]["evidence_observed_at_unix_ms"] = json!(now); + value["implementation"]["evidence_valid_for_ms"] = json!(600_000); + value["implementation"]["window_cost_model"]["cost"]["observed_at_unix_ms"] = json!(now); + value["implementation"]["window_cost_model"]["cost"]["valid_for_ms"] = json!(600_000); + value["environment"]["observed_at_unix_ms"] = json!(now); + value["environment"]["activation_unix_ms"] = json!(now); + value["environment"]["max_evidence_age_ms"] = json!(600_000); + value["environment"]["capability_snapshot_id"] = json!("promql-compliance"); + value["environment"]["dataset_identity"] = json!({ + "namespace": "promql-compliance", "dataset": dataset.name, + }); + // Both paths describe the actual finite replay. Differential data can + // have irregular gaps; its cadence contract bounds the largest gap. + // Benefit requires uniform cadence for comparable maintenance demand. + let (rate, cadence, count) = if benefit { + dataset.uniform_demand()? + } else { + dataset.replay_demand()? + }; + value["data_workload"]["ingestion_rate"]["value"] = json!(rate); + for (name, n) in [ + ("input_cardinality", dataset.series.len() as u64), + ("ingestion_volume", count as u64), + ("data_ingestion_interval", cadence), + ] { + value["data_workload"][name] = + json!({"value":n,"source":"declared","observed_at_ms":null,"valid_for_ms":null}); + } + value["implementation"]["scrape_interval_ms"] = json!(cadence); + // This runner controls the complete finite replay population. Its bounds + // are a closed-world fixture contract, never an inference about live data. + let samples: Vec = dataset + .series + .iter() + .flat_map(|series| series.samples.iter().map(|s| s.value)) + .collect(); + let lower = samples + .iter() + .copied() + .reduce(f64::min) + .context("empty replay")?; + let upper = samples + .iter() + .copied() + .reduce(f64::max) + .context("empty replay")?; + value["implementation"]["data_snapshot_id"] = + json!(format!("finite-replay-{}-{now}", dataset.name)); + for query in &suite.queries { + let root = control_plane::query_parser::parse_query_expr_canonical( + &query.expr, + planner_types::types::AccuracyTarget::EpsilonDelta { + epsilon: 0.01, + delta: 0.01, + }, + ) + .map_err(|e| anyhow::anyhow!(e.to_string()))?; + let planner_types::pre_asap::QueryExpr::BinaryOp { lhs, rhs, .. } = root else { + continue; + }; + let domains: Vec<_> = [lhs, rhs].into_iter().filter(|operand| matches!(operand.as_ref(), + planner_types::pre_asap::QueryExpr::Aggregate { measures, .. } + if matches!(measures.as_slice(), [planner_types::pre_asap::AggIntent::Quantile { .. } | planner_types::pre_asap::AggIntent::Avg { .. }]) + )).map(|operand| json!({"operand":operand,"lower":lower,"upper":upper, + "max_samples":samples.len(),"contract":"complete finite fixture replay; no other writers"})).collect(); + if domains.is_empty() { + continue; + } + value["implementation"]["accuracy_evidence"][&query.expr] = json!({ + "query_string": query.expr, + "data_snapshot_id": value["implementation"]["data_snapshot_id"], + "data_workload": value["data_workload"], + "source":"complete-finite-replay", "observed_at_unix_ms":now,"valid_for_ms":600_000, + "quantile_operand_domains": domains + }); + } + Ok(serde_json::from_value(value)?) +} + +pub fn validate_cost(plan: &CompiledPhysicalPlan) -> Result<()> { + let report = plan + .cost_comparison + .as_ref() + .context("missing cost comparison")?; + ensure!( + report.model_version == "backend-workload-resources-v2", + "automatic cost path was bypassed" + ); + ensure!( + report.selected_plan_id == plan.envelope.plan_id + && report.selected_manifest.plan_id == report.selected_plan_id, + "selected identity mismatch" + ); + ensure!( + !report.component_costs.is_empty() + && report + .component_costs + .keys() + .eq(report.selected_manifest.components.keys()), + "incomplete selected cost coverage" + ); + let mut minimum = f64::INFINITY; + let mut winner = None; + for candidate in &report.candidate_evaluations { + let selected = candidate.status == CandidateEvaluationStatus::Selected; + let Some(total) = candidate.total_cost else { + ensure!(!selected, "selected uncosted candidate"); + continue; + }; + let automatic = candidate + .automatic_cost + .as_ref() + .context("priced candidate has no resource breakdown")?; + ensure!( + automatic.model_version == report.model_version && !automatic.components.is_empty(), + "invalid model/empty components" + ); + ensure!( + automatic.weights["cpu_seconds"] == 1. + && automatic.weights["memory_byte_seconds"] == 1e-9 + && automatic.weights["network_bytes"] == 1e-8, + "unexpected resource weights" + ); + let mut sum = 0.; + for (id, r) in &automatic.components { + ensure!( + [r.cpu_seconds, r.memory_byte_seconds, r.network_bytes] + .iter() + .all(|n| n.is_finite() && *n >= 0.), + "invalid resource quantity {id}" + ); + ensure!( + (r.source == "analytical" && r.erp_record_ids.is_empty()) + || (r.source == "erp+analytical" && !r.erp_record_ids.is_empty()), + "invalid ERP provenance {id}" + ); + let cost = r.cpu_seconds + r.memory_byte_seconds * 1e-9 + r.network_bytes * 1e-8; + ensure!(cost.is_finite(), "resource cost overflow"); + sum += cost; + if selected { + ensure!( + report + .component_costs + .get(id) + .is_some_and(|n| equal(*n, cost)), + "selected component mismatch {id}" + ); + } + } + ensure!( + sum.is_finite() && total.is_finite() && equal(sum, total), + "inconsistent candidate total" + ); + minimum = minimum.min(total); + if selected { + ensure!( + winner.is_none() + && candidate.plan_id == Some(report.selected_plan_id) + && automatic.components.len() == report.component_costs.len(), + "selected candidate identity/coverage mismatch" + ); + winner = Some(total); + } + } + ensure!( + winner.is_some_and(|n| equal(n, minimum)), + "winner is not the least fully costed candidate" + ); + Ok(()) +} +fn equal(a: f64, b: f64) -> bool { + (a - b).abs() <= 1e-10 * a.abs().max(b.abs()).max(1.) +} + +pub fn validate_local(plan: &CompiledPhysicalPlan) -> Result<()> { + ensure!(!plan.query_plan.entries.is_empty(), "no installed queries"); + for entry in plan.query_plan.entries.values() { + ensure!(!entry.nodes.is_empty(), "query has no nodes"); + for node in entry.nodes.values() { + ensure!( + !matches!( + node, + QueryPlanNode::ExactFallback { .. } + | QueryPlanNode::ExternalExact { .. } + | QueryPlanNode::Logical { + operator: ResidualQueryOperator::ExactSubquery { .. } + | ResidualQueryOperator::CandidateExactSubquery { .. }, + .. + } + ), + "{} requires external exact execution", + entry.query_id + ); + } + } + Ok(()) +} + +fn gcd(mut a: u64, mut b: u64) -> u64 { + while b != 0 { + (a, b) = (b, a % b); + } + a +} + +/// Publish the exact costed generation; serving must not run candidate selection again. +pub fn installation( + plan: CompiledPhysicalPlan, +) -> asap_types::plan_publication::PhysicalPlanInstallRequest { + asap_types::plan_publication::PhysicalPlanInstallRequest { + summary_catalog: plan.summary_catalog, + collector_plans: plan.collector_plans, + precompute_plan: plan.precompute_plan, + transmission_plan: plan.transmission_plan, + query_plan: plan.query_plan, + // Match the ASAPQuery profile's local routing, as snapshot startup did. + storage_routing: None, + adaptation_evidence: Vec::new(), + } +} + +pub const FIXTURE_COST_MODEL: &str = "synthetic-execution-fixture-v1"; + +/// Price the exposed inventory with explicit test numbers, never online ERP. +/// Local execution is a fixture feasibility constraint, not a semantic rewrite. +pub fn with_fixture_costs( + mut input: BackendLocalPlanningInput, +) -> Result { + use control_plane::physical::{ + compiler::{DeploymentPlanCompiler, BACKEND_REVISION, PLANNER_REVISION}, + workload_cost::{ + enumerate_exact_and_materialized_candidates, manifest, WorkloadCostEvidence, + WorkloadQuote, + }, + }; + ensure!( + input.physical_inputs.erp.is_none(), + "execution fixture must not consume ERP" + ); + let (request, env) = input.clone().into_physical_compilation_request()?; + let mut seen = std::collections::BTreeSet::new(); + let mut quotes = Vec::new(); + for candidate in enumerate_exact_and_materialized_candidates(request)? { + let Ok(plan) = DeploymentPlanCompiler.compile_promql(candidate.clone(), env.clone()) else { + continue; // Production selection retains the binding rejection. + }; + let manifest = manifest(&plan, &candidate.queries)?; + if !seen.insert(serde_json::to_string(&manifest)?) { + continue; + } + quotes.push(WorkloadQuote { + executable: validate_local(&plan).is_ok(), + unit_costs: manifest + .components + .keys() + .map(|key| (key.clone(), 1.0)) + .collect(), + manifest, + }); + } + ensure!(!quotes.is_empty(), "fixture has no priceable candidates"); + input.workload_cost_evidence = Some(WorkloadCostEvidence { + backend_revision: BACKEND_REVISION.into(), + planner_revision: PLANNER_REVISION.into(), + data_snapshot_id: input + .physical_inputs + .data_snapshot_id + .clone() + .context("fixture data snapshot identity missing")?, + model_version: FIXTURE_COST_MODEL.into(), + observed_at_unix_ms: env.observed_at_unix_ms, + valid_for_ms: env.max_evidence_age_ms, + quotes, + }); + Ok(input) +} + +pub fn validate_fixture_cost(plan: &CompiledPhysicalPlan) -> Result<()> { + let report = plan + .cost_comparison + .as_ref() + .context("missing fixture cost comparison")?; + ensure!( + report.model_version == FIXTURE_COST_MODEL, + "fixture costs were bypassed" + ); + ensure!( + report.selected_plan_id == plan.envelope.plan_id + && report.selected_manifest.plan_id == plan.envelope.plan_id, + "selected identity mismatch" + ); + ensure!( + !report.component_costs.is_empty() + && report + .component_costs + .keys() + .eq(report.selected_manifest.components.keys()), + "incomplete fixture cost coverage" + ); + ensure!( + report + .component_costs + .values() + .all(|v| v.is_finite() && *v >= 0.), + "invalid fixture price" + ); + let mut winner = None; + let mut minimum = f64::INFINITY; + for candidate in &report.candidate_evaluations { + ensure!( + candidate.automatic_cost.is_none(), + "fixture acquired measured/analytical resource claims" + ); + if let Some(total) = candidate.total_cost { + ensure!(total.is_finite() && total >= 0., "invalid candidate price"); + minimum = minimum.min(total); + if candidate.status == CandidateEvaluationStatus::Selected { + ensure!( + winner.is_none() && candidate.plan_id == Some(plan.envelope.plan_id), + "invalid winner identity" + ); + winner = Some(total); + } + } else { + ensure!( + candidate.status != CandidateEvaluationStatus::Selected, + "selected unpriced candidate" + ); + } + } + let sum: f64 = report.component_costs.values().sum(); + ensure!( + winner.is_some_and(|cost| equal(cost, sum) && equal(cost, minimum)), + "inconsistent fixture selection" + ); + Ok(()) +} + +/// One target per distinct admitted deployment manifest. Binding/accuracy +/// failures remain in the inventory report and are never silently executable. +pub fn executable_fixture_targets( + input: &BackendLocalPlanningInput, + report: &control_plane::physical::workload_cost::CandidatePlanSelectionReport, +) -> Result> { + let evidence = input + .workload_cost_evidence + .as_ref() + .context("missing fixture quotes")?; + let admitted: std::collections::BTreeSet<_> = report + .candidate_evaluations + .iter() + .filter(|c| { + matches!( + c.status, + CandidateEvaluationStatus::Selected | CandidateEvaluationStatus::Unselected + ) + }) + .map(|c| c.plan_id.context("admitted candidate lacks identity")) + .collect::>()?; + let targets: Vec<_> = evidence + .quotes + .iter() + .filter(|q| q.executable && admitted.contains(&q.manifest.plan_id)) + .map(|q| q.manifest.clone()) + .collect(); + ensure!(!targets.is_empty(), "no executable fixture candidates"); + ensure!( + targets + .iter() + .map(|m| m.plan_id) + .collect::>() + == admitted, + "admitted candidate lacks an executable quote" + ); + Ok(targets) +} + +pub fn prefer_fixture_candidate( + mut input: BackendLocalPlanningInput, + target: &control_plane::physical::workload_cost::WorkloadCostManifest, +) -> Result { + let evidence = input + .workload_cost_evidence + .as_mut() + .context("missing fixture quotes")?; + ensure!( + evidence.model_version == FIXTURE_COST_MODEL, + "not fixture costs" + ); + let mut matches = 0; + for quote in &mut evidence.quotes { + let preferred = "e.manifest == target; + if preferred { + ensure!(quote.executable, "target is infeasible"); + matches += 1; + } + for value in quote.unit_costs.values_mut() { + *value = if preferred { 1.0 } else { 1e12 }; + } + } + ensure!(matches == 1, "target must identify exactly one quote"); + Ok(input) +} diff --git a/promql-compliance/runner/src/runner.rs b/promql-compliance/runner/src/runner.rs new file mode 100644 index 000000000..4eac85069 --- /dev/null +++ b/promql-compliance/runner/src/runner.rs @@ -0,0 +1,668 @@ +use crate::{ + compare, + compose::Compose, + input::{self, Dataset, Policy, Suite, Tolerance}, + planning, sql, transport, +}; +use anyhow::{ensure, Context, Result}; +use clap::Parser; +use reqwest::Client; +use serde_json::{json, Value}; +use std::{ + collections::BTreeMap, + path::{Path, PathBuf}, + time::{Instant, SystemTime, UNIX_EPOCH}, +}; + +#[derive(Parser, Debug, Clone)] +pub struct Args { + #[arg(long)] + pub dataset: PathBuf, + #[arg(long)] + pub suite: PathBuf, + #[arg( + long = "reference-url", + alias = "prometheus-url", + default_value = "http://127.0.0.1:19090" + )] + pub reference: String, + #[arg( + long = "test-url", + alias = "backend-url", + default_value = "http://127.0.0.1:19091" + )] + pub backend: String, + #[arg(long, default_value = "http://127.0.0.1:18428")] + pub victoria_url: String, + #[arg(long, default_value = "http://127.0.0.1:18123")] + pub clickhouse_url: String, + #[arg(long, default_value = "compliance-report.json")] + pub output: PathBuf, + #[arg(long)] + pub compose_file: Vec, + #[arg(long, default_value = "asapquery-rust-compliance")] + pub compose_project: String, + #[arg(long, default_value = "compliance-logs")] + pub logs_dir: PathBuf, + #[arg(long)] + pub keep_services: bool, + #[arg(long)] + pub base_time_ms: Option, + #[arg(long, default_value_t = 3)] + pub warmups: usize, + #[arg(long, default_value_t = 10)] + pub trials: usize, + /// Mutate fixture quotes and execute every admitted candidate for each query and the full ensemble. + #[arg(long)] + pub all_candidates: bool, +} +pub fn write_json(path: &Path, value: &impl serde::Serialize) -> Result<()> { + if let Some(parent) = path.parent().filter(|p| !p.as_os_str().is_empty()) { + std::fs::create_dir_all(parent)?; + } + std::fs::write(path, serde_json::to_vec_pretty(value)?)?; + Ok(()) +} +pub async fn run(args: Args, benefit: bool) -> Result<()> { + let result = if args.all_candidates { + ensure!( + !benefit && !args.keep_services && !args.compose_file.is_empty(), + "candidate sweep requires isolated Compose execution without benefit/keep-services" + ); + execute_candidates(&args).await + } else { + execute(&args, benefit).await + }; + match result { + Ok(report) => { + write_json(&args.output, &report)?; + ensure!( + report["passed"] == true, + "acceptance failed; inspect {}", + args.output.display() + ); + Ok(()) + } + Err(e) => { + write_json( + &args.output, + &json!({"passed":false,"benefitPassed":false,"dataset":args.dataset,"suite":args.suite,"error":format!("{e:#}")}), + )?; + Err(e) + } + } +} +async fn execute(args: &Args, benefit: bool) -> Result { + let data = Dataset::load(&args.dataset)?; + let suite = Suite::load(&args.suite)?; + if benefit { + ensure!( + suite.queries.len() == 10 + && suite + .queries + .iter() + .all(|q| !q.instant_offsets_seconds.is_empty()) + && args.trials >= 2 + && !args.compose_file.is_empty(), + "benefit needs ten shared instant cases, two trials and Compose" + ); + } + let now = SystemTime::now().duration_since(UNIX_EPOCH)?.as_millis() as u64; + let base = args.base_time_ms.unwrap_or(now as i64 - 1_800_000); + let snapshot = planning::snapshot(&suite, &data, now, base, benefit) + .context("compile unquoted workload snapshot")?; + let snapshot = if benefit { + snapshot + } else { + planning::with_fixture_costs(snapshot)? + }; + let snapshot_path = args.output.with_extension("snapshot.json"); + write_json(&snapshot_path, &snapshot)?; + // Compile before starting services. Preserve real planning failures without + // spending a container build or substituting a hand-selected candidate. + let plan = snapshot + .compile_promql() + .context("compile unquoted workload snapshot")?; + let plan_path = args.output.with_extension("plan.json"); + write_json(&plan_path, &plan)?; + if benefit { + planning::validate_cost(&plan).context("historical automatic cost gate")?; + } else { + planning::validate_fixture_cost(&plan).context("synthetic workload cost gate")?; + } + planning::validate_local(&plan).context("ASAP-local plan gate")?; + execute_installed(args, benefit, &suite, &data, base, plan).await +} + +async fn execute_installed( + args: &Args, + benefit: bool, + suite: &Suite, + data: &Dataset, + base: i64, + plan: control_plane::physical::compiler::CompiledPhysicalPlan, +) -> Result { + let plan_path = args.output.with_extension("plan.json"); + let installation_path = args.output.with_extension("install.json"); + write_json(&installation_path, &planning::installation(plan))?; + let installation_path = std::fs::canonicalize(installation_path)?; + let mut compose = Compose::new( + args.compose_file.clone(), + args.compose_project.clone(), + installation_path, + args.logs_dir.clone(), + args.keep_services, + ); + compose.start(benefit)?; + let client = transport::client()?; + for url in [ + format!("{}/api/v1/status/runtimeinfo", args.reference), + format!("{}/api/v1/health", args.backend), + ] { + transport::wait(&client, &url).await?; + } + let body = transport::encode(data, base)?; + let mut targets = vec![args.reference.as_str(), args.backend.as_str()]; + if benefit { + transport::wait(&client, &format!("{}/health", args.victoria_url)).await?; + transport::wait(&client, &format!("{}/ping", args.clickhouse_url)).await?; + targets.push(&args.victoria_url); + } + transport::push(&client, &body, &targets).await?; + transport::drain(&client, &args.backend).await?; + let semantic = compare_suite( + &client, + &args.reference, + &args.backend, + suite, + &data.name, + base, + ) + .await; + if !benefit { + compose.finish()?; + return Ok(semantic); + } + write_json(&args.output.with_extension("semantic.json"), &semantic)?; + ensure!( + semantic["passed"] == true, + "semantic gate failed; inspect semantic report" + ); + sql::seed(&client, &args.clickhouse_url, data, base).await?; + verify_baselines(&client, args, suite, data, base).await?; + let report = measure(&client, args, suite, data, &compose, base, &plan_path).await?; + compose.finish()?; + Ok(report) +} +fn response_result(result: Result) -> Value { + result.unwrap_or_else(|e| json!({"status":"error","error":format!("{e:#}")})) +} +fn backend_compare(reference: &Value, backend: &Value, policy: &Policy) -> Result<()> { + compare::compare(reference, backend, policy)?; + ensure!( + backend["servedBy"] + .as_str() + .is_some_and(|s| s == "asap_query"), + "backend response lacks ASAP provenance" + ); + Ok(()) +} +pub async fn compare_suite( + client: &Client, + reference: &str, + backend: &str, + suite: &Suite, + dataset: &str, + base: i64, +) -> Value { + let mut results = vec![]; + let mut passed = true; + for q in &suite.queries { + let policy = q.policy(&suite.comparison_defaults); + let mut report = json!({"name":q.name,"expr":q.expr,"tolerance":policy,"passed":true,"instant":[],"referenceParity":[],"backendParity":[]}); + let ranges = if let Some(r) = &q.range { + let a = response_result( + transport::query(client, reference, &q.expr, base, Some(r), false).await, + ); + let b = response_result( + transport::query(client, backend, &q.expr, base, Some(r), true).await, + ); + let outcome = compare::outcome(backend_compare(&a, &b, &policy)); + report["passed"] = outcome["passed"].clone(); + report["range"] = outcome; + report["rangeResponses"] = json!({"reference":a,"backend":b}); + Some((a, b)) + } else { + None + }; + for offset in &q.instant_offsets_seconds { + let at = match input::at_ms(base, *offset) { + Ok(at) => at, + Err(e) => { + report["passed"] = json!(false); + report["error"] = json!(e.to_string()); + continue; + } + }; + let a = response_result( + transport::query(client, reference, &q.expr, at, None, false).await, + ); + let b = + response_result(transport::query(client, backend, &q.expr, at, None, true).await); + let outcome = compare::outcome(backend_compare(&a, &b, &policy)); + report["passed"] = json!(report["passed"] == true && outcome["passed"] == true); + report["instant"].as_array_mut().unwrap().push(json!({"offsetSeconds":offset,"timeMs":at,"comparison":outcome,"responses":{"reference":a,"backend":b}})); + if let Some((ra, rb)) = &ranges { + for (name, range, instant) in + [("referenceParity", ra, &a), ("backendParity", rb, &b)] + { + let p = compare::outcome(compare::parity(range, instant, at, &policy)); + report["passed"] = json!(report["passed"] == true && p["passed"] == true); + report[name] + .as_array_mut() + .unwrap() + .push(json!({"offsetSeconds":offset,"timeMs":at,"comparison":p})); + } + } + } + passed &= report["passed"] == true; + results.push(report); + } + json!({"suite":suite.name,"dataset":dataset,"baseTimeMs":base,"queries":results,"passed":passed}) +} +async fn verify_baselines( + client: &Client, + args: &Args, + suite: &Suite, + data: &Dataset, + base: i64, +) -> Result<()> { + let exact = Policy { + value_tolerance: Some(Tolerance { + absolute: Some(1e-8), + relative: Some(1e-8), + }), + }; + for q in &suite.queries { + let at = input::at_ms(base, q.instant_offsets_seconds[0])?; + let prom = transport::query(client, &args.reference, &q.expr, at, None, false).await?; + let vm = transport::query(client, &args.victoria_url, &q.expr, at, None, false).await?; + compare::compare(&prom, &vm, &q.policy(&suite.comparison_defaults)) + .with_context(|| format!("VictoriaMetrics {}", q.name))?; + let sql = sql::baseline(&q.name, at, sql::window_ms(&q.expr)?)?; + let rows = sql::rows(client, &args.clickhouse_url, &sql).await?; + let mut values = vec![]; + for row in rows { + let labels = if let Some(id) = row["series_id"] + .as_u64() + .or_else(|| row["series_id"].as_str().and_then(|s| s.parse().ok())) + { + data.series + .get(id.checked_sub(1).context("invalid SQL series id")? as usize) + .context("unknown SQL series id")? + .labels + .clone() + } else { + BTreeMap::from([( + "label_0".into(), + row["label_0"] + .as_str() + .context("SQL row has no group")? + .to_owned(), + )]) + }; + let number = row["value"].as_f64().context("non-numeric SQL value")?; + values.push(json!({"metric":labels,"value":[at as f64/1000.,number.to_string()]})); + } + let mut prom = prom; + // SQL series IDs resolve to complete source labels; metric-name + // retention differs between native SQL and PromQL functions. + for row in prom["data"]["result"] + .as_array_mut() + .context("expected PromQL vector")? + { + row["metric"] + .as_object_mut() + .context("invalid labels")? + .remove("__name__"); + } + let sql_response = + json!({"status":"success","data":{"resultType":"vector","result":values}}); + compare::compare(&prom, &sql_response, &exact) + .with_context(|| format!("ClickHouse {}", q.name))?; + } + Ok(()) +} +async fn execute_query( + client: &Client, + args: &Args, + name: &str, + q: &input::Query, + at: i64, +) -> Result<()> { + if name == "clickhouse" { + sql::rows( + client, + &args.clickhouse_url, + &sql::baseline(&q.name, at, sql::window_ms(&q.expr)?)?, + ) + .await?; + } else { + let url = match name { + "backend" => &args.backend, + "victoria" => &args.victoria_url, + _ => &args.reference, + }; + let response = transport::query(client, url, &q.expr, at, None, name == "backend").await?; + ensure!( + response["status"] == "success", + "measurement query failed: {response}" + ); + if name == "backend" { + ensure!( + response["servedBy"] + .as_str() + .is_some_and(|s| s == "asap_query"), + "fallback during performance measurement" + ); + } + } + Ok(()) +} +pub fn percentile(samples: &[f64], p: f64) -> Result { + ensure!( + !samples.is_empty() + && samples.iter().all(|n| n.is_finite() && *n >= 0.) + && p > 0. + && p <= 1., + "invalid latency samples" + ); + let mut sorted = samples.to_vec(); + sorted.sort_by(f64::total_cmp); + Ok(sorted[(p * sorted.len() as f64).ceil() as usize - 1]) +} +async fn measure( + client: &Client, + args: &Args, + suite: &Suite, + data: &Dataset, + compose: &Compose, + base: i64, + plan: &Path, +) -> Result { + let mut targets = serde_json::Map::new(); + for (name, service) in [ + ("backend", "data-plane"), + ("prometheus", "prometheus"), + ("victoria", "victoria"), + ("clickhouse", "clickhouse"), + ] { + for q in &suite.queries { + for _ in 0..args.warmups { + execute_query( + client, + args, + name, + q, + input::at_ms(base, q.instant_offsets_seconds[0])?, + ) + .await?; + } + } + let before = compose.usage(service)?; + let mut queries = serde_json::Map::new(); + for q in &suite.queries { + let mut samples = vec![]; + let at = input::at_ms(base, q.instant_offsets_seconds[0])?; + for _ in 0..args.trials { + let start = Instant::now(); + execute_query(client, args, name, q, at).await?; + samples.push(start.elapsed().as_secs_f64() * 1000.); + } + queries.insert(q.name.clone(),json!({"p50Ms":percentile(&samples,0.5)?,"p95Ms":percentile(&samples,0.95)?,"samplesMs":samples})); + } + let after = compose.usage(service)?; + let cpu = after + .0 + .checked_sub(before.0) + .context("CPU counter reset during measurement")?; + targets.insert( + name.into(), + json!({"queries":queries,"cpuUsec":cpu,"memoryPeakBytes":after.1}), + ); + } + let failures = benefit_failures(&Value::Object(targets.clone()), suite)?; + Ok( + json!({"suite":suite.name,"dataset":data.name,"baseTimeMs":base,"selectedPlan":plan,"warmups":args.warmups,"trials":args.trials,"targets":targets,"failures":failures,"passed":failures.is_empty(),"benefitPassed":failures.is_empty()}), + ) +} +pub fn benefit_failures(targets: &Value, suite: &Suite) -> Result> { + let mut failures = vec![]; + for name in ["prometheus", "victoria", "clickhouse"] { + for metric in ["cpuUsec", "memoryPeakBytes"] { + let backend = targets["backend"][metric] + .as_u64() + .context("missing backend usage")?; + let baseline = targets[name][metric] + .as_u64() + .context("missing baseline usage")?; + if backend >= baseline { + failures.push(format!("backend {metric} >= {name}")); + } + } + for q in &suite.queries { + let backend = targets["backend"]["queries"][&q.name]["p95Ms"] + .as_f64() + .context("missing backend latency")?; + let baseline = targets[name]["queries"][&q.name]["p95Ms"] + .as_f64() + .context("missing baseline latency")?; + ensure!( + backend.is_finite() && baseline.is_finite() && backend >= 0. && baseline >= 0., + "invalid latency" + ); + if backend >= baseline { + failures.push(format!("{} backend p95 >= {name}", q.name)); + } + } + } + Ok(failures) +} +pub fn report_card(directory: &Path) -> Result<()> { + let mut cases = vec![]; + let mut passed = true; + for entry in std::fs::read_dir(directory)? { + let path = entry?.path(); + let filename = path.file_name().unwrap().to_string_lossy(); + if path.extension().is_none_or(|s| s != "json") + || filename == "summary.json" + || [ + ".snapshot.json", + ".plan.json", + ".install.json", + ".semantic.json", + ] + .iter() + .any(|suffix| filename.ends_with(suffix)) + { + continue; + } + let report: Value = serde_json::from_slice(&std::fs::read(path)?)?; + ensure!( + report.get("passed").and_then(Value::as_bool).is_some(), + "not a compliance report" + ); + passed &= report["passed"] == true; + let queries = report["queries"].as_array().cloned().unwrap_or_default(); + let (mut local, mut fallback) = (0, 0); + for q in &queries { + let responses = std::iter::once(&q["rangeResponses"]["backend"]).chain( + q["instant"] + .as_array() + .into_iter() + .flatten() + .map(|i| &i["responses"]["backend"]), + ); + for response in responses.filter(|r| !r.is_null()) { + if response["servedBy"] + .as_str() + .is_some_and(|s| s == "asap_query") + { + local += 1 + } else { + fallback += 1 + } + } + } + cases.push(json!({"dataset":report["dataset"],"suite":report["suite"],"passed":report["passed"],"queries":queries.len(),"asapQuery":local,"prometheusFallback":fallback})); + } + ensure!(!cases.is_empty(), "no comparison reports"); + cases.sort_by_key(|c| c["dataset"].to_string()); + let mut markdown=format!("# PromQL compliance report card\n\nOverall: **{passed}**\n\n| Dataset | Suite | Passed | Queries | ASAPQuery | Fallback |\n|---|---|---|---:|---:|---:|\n"); + for c in &cases { + markdown.push_str(&format!( + "| {} | {} | {} | {} | {} | {} |\n", + c["dataset"], + c["suite"], + c["passed"], + c["queries"], + c["asapQuery"], + c["prometheusFallback"] + )); + } + write_json( + &directory.join("summary.json"), + &json!({"cases":cases,"passed":passed}), + )?; + std::fs::write(directory.join("summary.md"), markdown)?; + Ok(()) +} + +/// Exercise each admitted physical workload by changing quotes, never by +/// replacing the selected DAG after the production selector has run. +async fn execute_candidates(args: &Args) -> Result { + let data = Dataset::load(&args.dataset)?; + let suite = Suite::load(&args.suite)?; + let mut scopes: Vec<(String, Suite)> = suite + .queries + .iter() + .map(|query| { + ( + format!("query:{}", query.name), + Suite { + queries: vec![query.clone()], + ..suite.clone() + }, + ) + }) + .collect(); + for (name, names) in [ + ( + "shared-rate", + vec!["temporal-rate", "grouped-rate", "topk-rate"], + ), + ( + "shared-quantiles", + vec!["temporal-quantile", "quantile-ratio"], + ), + ] { + let queries: Vec<_> = suite + .queries + .iter() + .filter(|q| names.contains(&q.name.as_str())) + .cloned() + .collect(); + if queries.len() == names.len() { + scopes.push(( + name.into(), + Suite { + queries, + ..suite.clone() + }, + )); + } + } + if suite.queries.len() > 1 { + scopes.push(("full-ensemble".into(), suite)); + } + let mut runs = Vec::new(); + let mut inventories = Vec::new(); + let mut results = Vec::new(); + let mut passed = true; + for (scope_index, (name, suite)) in scopes.into_iter().enumerate() { + let now = SystemTime::now().duration_since(UNIX_EPOCH)?.as_millis() as u64; + let base = args.base_time_ms.unwrap_or(now as i64 - 1_800_000); + let input = + planning::with_fixture_costs(planning::snapshot(&suite, &data, now, base, false)?)?; + let baseline = input.clone().compile_promql()?; + planning::validate_fixture_cost(&baseline)?; + let report = baseline.cost_comparison.context("missing inventory")?; + let targets = planning::executable_fixture_targets(&input, &report)?; + inventories.push(json!({"scope":name,"queries":suite.queries.iter().map(|q| &q.expr).collect::>(), + "candidateEvaluations": report.candidate_evaluations,"executableCandidates":targets.len(), + "searchScope": report.materialization_search_coverage})); + for (candidate_index, target) in targets.into_iter().enumerate() { + let mut trial = args.clone(); + let directory = args + .output + .with_extension("") + .join(format!("scope-{scope_index}")); + trial.output = directory.join(format!("candidate-{candidate_index}.json")); + trial.logs_dir = args + .logs_dir + .join(format!("scope-{scope_index}-candidate-{candidate_index}")); + trial.compose_project = + format!("{}-s{scope_index}-c{candidate_index}", args.compose_project); + eprintln!( + "Executing {name} candidate {candidate_index}: plan {}", + target.plan_id + ); + let trial_result: Result = async { + let priced = planning::prefer_fixture_candidate(input.clone(), &target)?; + write_json(&trial.output.with_extension("snapshot.json"), &priced)?; + let selected = priced.compile_promql()?; + planning::validate_fixture_cost(&selected)?; + planning::validate_local(&selected)?; + ensure!( + selected + .cost_comparison + .as_ref() + .context("missing selected report")? + .selected_manifest + == target, + "synthetic prices did not select the intended candidate" + ); + write_json(&trial.output.with_extension("plan.json"), &selected)?; + execute_installed(&trial, false, &suite, &data, base, selected).await + } + .await; + let result = match trial_result { + Ok(result) => result, + Err(error) => json!({"passed":false,"error":format!("{error:#}")}), + }; + passed &= result["passed"] == true; + write_json(&trial.output, &result)?; + if let Some(queries) = result["queries"].as_array() { + results.extend(queries.iter().cloned().map(|mut query| { + query["candidatePlanId"] = json!(target.plan_id); + query["scope"] = json!(name); + query + })); + } + runs.push( + json!({"scope":name,"candidatePlanId":target.plan_id,"passed":result["passed"], + "report":trial.output,"error":result.get("error")}), + ); + // Persist partial coverage so an interrupted run never looks complete. + write_json( + &args.output, + &json!({"passed":false,"complete":false,"inventories":inventories,"candidateRuns":runs}), + )?; + } + } + ensure!(!runs.is_empty(), "no executable candidate trials"); + Ok( + json!({"passed":passed,"complete":true,"dataset":data.name,"inventories":inventories, + "candidateRuns":runs,"queries":results,"costModel":planning::FIXTURE_COST_MODEL}), + ) +} diff --git a/promql-compliance/runner/src/sql.rs b/promql-compliance/runner/src/sql.rs new file mode 100644 index 000000000..68ece349f --- /dev/null +++ b/promql-compliance/runner/src/sql.rs @@ -0,0 +1,115 @@ +use crate::input::{at_ms, Dataset}; +use anyhow::{bail, ensure, Context, Result}; +use reqwest::Client; +use serde_json::{json, Value}; + +pub fn window_ms(expr: &str) -> Result { + let regex = regex::Regex::new(r"\[(\d+)([smhd])\]")?; + let Some(c) = regex.captures(expr) else { + return Ok(60_000); + }; + let n: i64 = c[1].parse()?; + let unit = match &c[2] { + "s" => 1000, + "m" => 60_000, + "h" => 3_600_000, + _ => 86_400_000, + }; + n.checked_mul(unit).context("window overflow") +} +pub fn baseline(name: &str, evaluation_ms: i64, window_ms: i64) -> Result { + ensure!(window_ms > 0, "positive window required"); + let (counter, query) = match name { + "spatial-sum" => ( + false, + r#"SELECT label_0, sum(value) AS value FROM instant_samples GROUP BY label_0"#, + ), + "spatial-topk" => ( + false, + r#"SELECT series_id, label_0, value FROM instant_samples ORDER BY label_0, value DESC, series_id LIMIT 3 BY label_0"#, + ), + "spatial-quantile" => ( + false, + r#"SELECT label_0, quantileExactInclusive(0.9)(value) AS value FROM instant_samples GROUP BY label_0"#, + ), + "temporal-sum" => ( + false, + r#"SELECT series_id, sum(value) AS value FROM window_samples GROUP BY series_id"#, + ), + "temporal-quantile" => ( + false, + r#"SELECT series_id, quantileExactInclusive(0.9)(value) AS value FROM window_samples GROUP BY series_id"#, + ), + "temporal-rate" => ( + true, + r#"SELECT series_id, rate_value AS value FROM per_series_counter"#, + ), + "grouped-rate" => ( + true, + r#"SELECT label_0, sum(rate_value) AS value FROM per_series_counter GROUP BY label_0"#, + ), + "grouped-temporal-sum" => ( + false, + r#"SELECT label_0, sum(series_sum) AS value FROM (SELECT series_id, label_0, sum(value) AS series_sum FROM window_samples GROUP BY series_id, label_0) GROUP BY label_0"#, + ), + "topk-rate" => ( + true, + r#"SELECT series_id, label_0, rate_value AS value FROM per_series_counter ORDER BY label_0, rate_value DESC, series_id LIMIT 3 BY label_0"#, + ), + "quantile-ratio" => ( + false, + r#"SELECT series_id, quantileExactInclusive(0.9)(value) / quantileExactInclusive(0.5)(value) AS value FROM window_samples GROUP BY series_id"#, + ), + _ => bail!("no ClickHouse baseline for {name}"), + }; + let prefix = include_str!("../sql/prefix.sql") + .replace("{evaluation_ms}", &evaluation_ms.to_string()) + .replace("{window_ms}", &window_ms.to_string()); + Ok(format!( + "{prefix}{}\n{query} FORMAT JSON", + if counter { + include_str!("../sql/counter.sql") + } else { + "" + } + )) +} +pub async fn post(client: &Client, url: &str, sql: &str, body: Option) -> Result { + let request = client.post(format!("{}/", url.trim_end_matches('/'))); + let request = if let Some(body) = body { + request.query(&[("query", sql)]).body(body) + } else { + request.body(sql.to_owned()) + }; + let response = request.send().await?; + let status = response.status(); + let text = response.text().await?; + ensure!(status.is_success(), "ClickHouse {status}: {text}"); + Ok(text) +} +pub async fn rows(client: &Client, url: &str, sql: &str) -> Result> { + let response: Value = serde_json::from_str(&post(client, url, sql, None).await?)?; + Ok(response["data"] + .as_array() + .context("missing ClickHouse rows")? + .clone()) +} +pub async fn seed(client: &Client, url: &str, data: &Dataset, base: i64) -> Result<()> { + post(client,url,"CREATE TABLE IF NOT EXISTS samples (series_id UInt64, label_0 String, ts_ms Int64, value Float64) ENGINE = MergeTree ORDER BY (series_id, ts_ms)",None).await?; + post(client, url, "TRUNCATE TABLE samples", None).await?; + let mut body = String::new(); + for (i, s) in data.series.iter().enumerate() { + for p in &s.samples { + body.push_str(&serde_json::to_string(&json!({"series_id":i+1,"label_0":s.labels.get("label_0").context("missing label_0")?,"ts_ms":at_ms(base,p.offset_seconds)?,"value":p.value}))?); + body.push('\n'); + } + } + post( + client, + url, + "INSERT INTO samples FORMAT JSONEachRow", + Some(body), + ) + .await?; + Ok(()) +} diff --git a/promql-compliance/runner/src/transport.rs b/promql-compliance/runner/src/transport.rs new file mode 100644 index 000000000..70b9feab4 --- /dev/null +++ b/promql-compliance/runner/src/transport.rs @@ -0,0 +1,167 @@ +use crate::input::{at_ms, Dataset, Range}; +use anyhow::{ensure, Context, Result}; +use prost::Message; +use reqwest::Client; +use serde_json::Value; +use std::time::Duration; +// Remote Write v1's wire messages. These tags are the same ones accepted by +// the backend and Prometheus; keep the minimal schema independently testable. +#[derive(Clone, PartialEq, Message)] +pub struct WriteRequest { + #[prost(message, repeated, tag = "1")] + pub timeseries: Vec, +} +#[derive(Clone, PartialEq, Message)] +pub struct TimeSeries { + #[prost(message, repeated, tag = "1")] + pub labels: Vec