From f489f27dc5071900d70b45d8f0070ee79eff1f62 Mon Sep 17 00:00:00 2001 From: zz_y Date: Mon, 28 Sep 2026 01:58:31 +0000 Subject: [PATCH] test: isolate observed workload statistics contracts --- control_plane/tests/workload_statistics.rs | 99 +++++++++++++++++++ .../workload-statistics-validation.md | 17 ++++ 2 files changed, 116 insertions(+) create mode 100644 control_plane/tests/workload_statistics.rs create mode 100644 docs/design_docs/workload-statistics-validation.md diff --git a/control_plane/tests/workload_statistics.rs b/control_plane/tests/workload_statistics.rs new file mode 100644 index 00000000..8bea9262 --- /dev/null +++ b/control_plane/tests/workload_statistics.rs @@ -0,0 +1,99 @@ +//! Statistics-consumer contracts. Controlled observations are not live telemetry. +use control_plane::physical::{compiler::BackendLocalPlanningInput, erp::ErpShapeObserver}; +use serde_json::{json, Value}; + +fn wire() -> Value { + let mut wire: Value = serde_json::from_str(include_str!( + "../../docs/examples/asapquery-planning-snapshot.json" + )) + .unwrap(); + wire["query_workload"]["repeating_queries"][0]["query"] = json!("sum_over_time(m[1m])"); + wire["data_workload"]["ingestion_rate"] = json!({ + "value":25.0,"source":"observed","observed_at_ms":9500,"valid_for_ms":1000 + }); + wire["data_workload"]["input_cardinality"] = json!({ + "value":125,"source":"observed","observed_at_ms":9500,"valid_for_ms":1000 + }); + wire +} + +/// Source rate, source cardinality and query cadence remain distinct quantities. +#[test] +fn observed_workload_facts_survive_binding_without_changing_units() { + let mut wire = wire(); + for interval in [5000, 20000] { + wire["query_workload"]["repeating_queries"][0]["demand"] = + json!({"fixed_interval_at":{"interval":interval,"evaluation_phase":0}}); + let input: BackendLocalPlanningInput = serde_json::from_value(wire.clone()).unwrap(); + let expected = input.data_workload.clone(); + let (request, _) = input.into_physical_compilation_request().unwrap(); + assert_eq!(request.data_workload.as_ref(), Some(&expected)); + assert_eq!( + request.queries[0] + .summary_lifecycle_inputs + .ingestion_rate_per_second, + 25.0 + ); + assert_eq!( + request.queries[0] + .summary_lifecycle_inputs + .evaluation_interval_ms, + interval + ); + assert_eq!(expected.input_cardinality.value, Some(125)); + } +} + +/// Missing, expired and future observations cannot be priced as zero ingestion. +#[test] +fn unavailable_rate_observations_are_rejected() { + for observation in [ + json!({"value":null,"source":"unknown","observed_at_ms":null,"valid_for_ms":null}), + json!({"value":25.0,"source":"observed","observed_at_ms":8000,"valid_for_ms":1000}), + json!({"value":25.0,"source":"observed","observed_at_ms":11000,"valid_for_ms":1000}), + ] { + let mut wire = wire(); + wire["data_workload"]["ingestion_rate"] = observation.clone(); + let input: BackendLocalPlanningInput = serde_json::from_value(wire).unwrap(); + assert!( + input.into_physical_compilation_request().is_err(), + "{observation}" + ); + } +} + +/// No recent arrivals does not imply no retained series. +#[test] +fn observed_zero_rate_does_not_erase_cardinality() { + let mut wire = wire(); + wire["data_workload"]["ingestion_rate"]["value"] = json!(0.0); + let input: BackendLocalPlanningInput = serde_json::from_value(wire).unwrap(); + let (request, _) = input.into_physical_compilation_request().unwrap(); + assert_eq!( + request.data_workload.unwrap().input_cardinality.value, + Some(125) + ); + assert_eq!( + request.queries[0] + .summary_lifecycle_inputs + .ingestion_rate_per_second, + 0.0 + ); +} + +/// Repeated events count toward workload, not distinct cardinality; overflow invalidates the window. +#[test] +fn observation_population_counts_events_and_distinct_keys_separately() { + let mut observer = ErpShapeObserver::with_limits(2, 4).unwrap(); + for _ in 0..30 { + observer.observe("series-a", 0).unwrap(); + } + for _ in 0..10 { + observer.observe("series-b", 1).unwrap(); + } + let observation = observer.snapshot().unwrap(); + assert_eq!(observation.observation.observed_events, 40); + assert_eq!(observation.observation.cardinality, 2); + assert!(observer.observe("series-c", 2).is_err()); + assert!(observer.snapshot().is_none()); +} diff --git a/docs/design_docs/workload-statistics-validation.md b/docs/design_docs/workload-statistics-validation.md new file mode 100644 index 00000000..6168bb3f --- /dev/null +++ b/docs/design_docs/workload-statistics-validation.md @@ -0,0 +1,17 @@ +# Workload statistics validation + +This test PR isolates the consumer contract from ranking and sketch accuracy. +`control_plane/tests/workload_statistics.rs` verifies units, retained-series +cardinality with zero arrivals, expired/future/missing observations, and bounded +population accounting. Its controlled observations are not live telemetry. + +The intended measurement sources remain remote_write accepted-sample counters, +query-tracker executions, and Prometheus series observations. Rates need explicit +observation windows and counter-reset handling; series scope must match the input +computation. Query cadence and ingestion cadence are different quantities. + +The existing discovery replay adapter reports a derived finite-replay rate and +declared query recurrence. Those are not live ingestion/query measurements. This +PR does not claim the three production collectors are fully wired together. +Level 3 must preserve original trace timing, observation provenance and scope; +missing observations make the real-evidence run incomplete.