From b832aa3e8f39de76f08787c14d5b8bd4bd1facfb Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 3 Sep 2026 17:48:36 -0600 Subject: [PATCH 1/2] feat(workload): serialize canonical planner snapshots --- crates/types/src/workload.rs | 127 ++++++++++++++++++++++++++++------- 1 file changed, 101 insertions(+), 26 deletions(-) diff --git a/crates/types/src/workload.rs b/crates/types/src/workload.rs index 512230c5..96ad2ca4 100644 --- a/crates/types/src/workload.rs +++ b/crates/types/src/workload.rs @@ -4,25 +4,30 @@ use serde::{Deserialize, Serialize}; // ── Query surface ───────────────────────────────────────────────────────────── /// A raw query string in its source language, before any parsing. -#[derive(Debug, Clone, PartialEq, Eq, Hash)] +#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)] +#[serde(transparent)] pub struct Query(pub String); /// How often a repeating query fires, in milliseconds. -#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] +#[serde(transparent)] pub struct RepetitionInterval(pub u32); /// Milliseconds since the Unix epoch. Workload timestamps use one explicit /// representation so schedules, observations, and time selections agree. -#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] +#[serde(transparent)] pub struct TimestampMs(pub u64); /// A non-negative duration in milliseconds. -#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] +#[serde(transparent)] pub struct DurationMs(pub u64); /// SQL dialect variant — different dialects have different syntax and /// function sets that affect how the query string is parsed. -#[derive(Debug, Clone, PartialEq, Eq, Hash)] +#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] pub enum SqlDialect { DataFusionSQL, ClickhouseSQL, @@ -31,7 +36,8 @@ pub enum SqlDialect { /// Source language of every query in the workload. /// All queries in a single `QueryWorkload` share the same language. -#[derive(Debug, Clone, PartialEq, Eq, Hash)] +#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] pub enum QueryLanguage { PromQL, SQL(SqlDialect), @@ -43,7 +49,8 @@ pub enum QueryLanguage { /// Whether the caller explicitly requested an accuracy target or inherited /// the normalized exact default. -#[derive(Debug, Clone, PartialEq)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] pub enum AccuracyRequirement { Explicit(AccuracyTarget), ImplicitExact, @@ -59,7 +66,8 @@ impl AccuracyRequirement { } /// Optional maximum wall-clock response time for one query execution. -#[derive(Debug, Clone, Copy, Default, PartialEq)] +#[derive(Debug, Clone, Copy, Default, PartialEq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] pub enum LatencyRequirement { ExplicitMaxMs(f64), #[default] @@ -68,7 +76,8 @@ pub enum LatencyRequirement { /// Independent accuracy and response-latency constraints attached to one /// query in the workload. -#[derive(Debug, Clone, PartialEq)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] pub struct QueryRequirements { pub accuracy: AccuracyRequirement, pub response_latency: LatencyRequirement, @@ -101,7 +110,8 @@ impl Default for QueryRequirements { /// }; /// let exploration = Predictability::AdHoc; /// ``` -#[derive(Debug, Clone, Default, PartialEq, Eq)] +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] pub enum Predictability { AdHoc, Predictable { @@ -129,7 +139,8 @@ pub enum Predictability { /// let historical_report = QueryTimeScope::Longitudinal; /// let week_over_week_dashboard = QueryTimeScope::Mixed; /// ``` -#[derive(Debug, Clone, Default, PartialEq, Eq)] +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] pub enum QueryTimeScope { RealTime, Longitudinal, @@ -156,7 +167,8 @@ pub enum QueryTimeScope { /// as_of: None, /// }; /// ``` -#[derive(Debug, Clone, Default, PartialEq, Eq)] +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] pub struct TimeSelection { pub scope: QueryTimeScope, pub lookback: Option, @@ -169,14 +181,16 @@ pub struct TimeSelection { /// /// For example, `Confidence(0.95)` says the demand estimate is supplied with /// 95% confidence. -#[derive(Debug, Clone, Copy, PartialEq)] +#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)] +#[serde(transparent)] pub struct Confidence(pub f64); /// Event-time interval from which a demand estimate was learned. /// /// For example, `{ start: TimestampMs(0), end: TimestampMs(60_000) }` /// describes an estimate based on the first minute of observations. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] pub struct ObservationWindow { pub start: TimestampMs, pub end: TimestampMs, @@ -212,7 +226,8 @@ pub struct ObservationWindow { /// valid_for: Some(DurationMs(5 * 60 * 1_000)), /// }; /// ``` -#[derive(Debug, Clone, PartialEq)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] pub struct DemandEstimate { pub observation_window: ObservationWindow, pub expected_rate: Rate, @@ -237,7 +252,8 @@ pub struct DemandEstimate { /// TimestampMs(200_000), /// ]); /// ``` -#[derive(Debug, Clone, PartialEq)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] pub enum RepeatedDemand { FixedInterval(RepetitionInterval), Scheduled(Vec), @@ -261,7 +277,8 @@ pub enum RepeatedDemand { /// RepeatedDemand::FixedInterval(RepetitionInterval(60_000)), /// ); /// ``` -#[derive(Debug, Clone, PartialEq)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] pub enum QueryRecurrence { OneTime { invocations: u64, @@ -294,7 +311,8 @@ pub enum QueryRecurrence { /// }, /// }; /// ``` -#[derive(Debug, Clone)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] pub struct BatchEntry { pub query: Query, pub requirements: QueryRequirements, @@ -325,7 +343,8 @@ pub struct BatchEntry { /// }, /// }; /// ``` -#[derive(Debug, Clone)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] pub struct RepeatingEntry { pub query: Query, pub demand: RepeatedDemand, @@ -357,7 +376,8 @@ pub struct RepeatingEntry { /// }, /// }; /// ``` -#[derive(Debug, Clone, PartialEq)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] pub struct QueryWorkloadEntry { pub query: Query, pub requirements: QueryRequirements, @@ -408,7 +428,8 @@ pub enum DataArrival { } /// Statistical distribution of keys in the incoming data stream. -#[derive(Debug, Clone, Default, PartialEq, Eq)] +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] pub enum DataDistribution { /// Zipf-distributed keys (s ≈ 1.1). A small number of keys dominate, /// so only a fraction of sketch cells are touched per window. Typical @@ -424,7 +445,8 @@ pub enum DataDistribution { } /// Where an empirical workload value came from. -#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] pub enum EvidenceSource { Declared, Observed, @@ -435,7 +457,8 @@ pub enum EvidenceSource { /// A workload value together with the provenance and freshness needed to /// decide whether it is safe to use. Times and durations are milliseconds. -#[derive(Debug, Clone, PartialEq)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] pub struct Evidence { pub value: Option, pub source: EvidenceSource, @@ -486,12 +509,14 @@ impl DemandEstimate { /// Queries per second, samples per second, or another rate whose unit is /// established by the field that contains it. -#[derive(Debug, Clone, Copy, PartialEq)] +#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)] +#[serde(transparent)] pub struct Rate(pub f64); /// Workload-level facts about the data being queried. Unlike the former /// ingestion-only `DataCharacteristics`, this also represents data at rest. -#[derive(Debug, Clone, Default, PartialEq)] +#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] pub struct DataWorkload { pub arrival: DataArrival, pub ingestion_volume: Evidence, @@ -508,7 +533,8 @@ pub struct DataWorkload { /// `query_batch` and `repeating_queries` may both be present. [`Self::entries`] /// normalizes them into one ordered stream without conflating recurrence with /// data arrival. -#[derive(Debug, Clone)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] pub struct QueryWorkload { /// Source language shared by all queries in this workload. pub language: QueryLanguage, @@ -732,6 +758,55 @@ mod tests { ); } + #[test] + fn canonical_workloads_have_a_strict_json_round_trip() { + let workload = QueryWorkload { + language: QueryLanguage::PromQL, + query_batch: None, + repeating_queries: Some(vec![RepeatingEntry { + query: Query("rate(requests_total[5m])".into()), + demand: RepeatedDemand::FixedInterval(RepetitionInterval(10_000)), + requirements: QueryRequirements { + accuracy: AccuracyRequirement::Explicit(AccuracyTarget::EpsilonDelta { + epsilon: 0.01, + delta: 0.001, + }), + response_latency: LatencyRequirement::ExplicitMaxMs(100.0), + }, + predictability: Predictability::Predictable { known_at: None }, + time_selection: TimeSelection { + scope: QueryTimeScope::RealTime, + lookback: Some(DurationMs(300_000)), + as_of: None, + }, + }]), + data_workload: Some(DataWorkload { + arrival: DataArrival::ContinuouslyIngesting, + ingestion_rate: Evidence { + value: Some(Rate(1_000.0)), + source: EvidenceSource::Declared, + observed_at_ms: None, + valid_for_ms: None, + }, + input_cardinality: Evidence { + value: Some(10_000), + source: EvidenceSource::Declared, + observed_at_ms: None, + valid_for_ms: None, + }, + ..Default::default() + }), + }; + + let json = serde_json::to_string_pretty(&workload).expect("serialize workload"); + let decoded: QueryWorkload = serde_json::from_str(&json).expect("deserialize workload"); + assert_eq!(decoded, workload); + assert!(json.contains("\"continuously_ingesting\"")); + + let with_unknown = json.replacen("{", "{\"unknown\":true,", 1); + assert!(serde_json::from_str::(&with_unknown).is_err()); + } + #[test] fn estimated_demand_validates_window_rate_and_confidence() { let mut workload = base_workload(); From 9949fcc997e1caca7811d97b7ee221b0110e7518 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 3 Sep 2026 19:28:35 -0600 Subject: [PATCH 2/2] fix(workload): stabilize strict snapshot wire names --- crates/types/src/types.rs | 1 + crates/types/src/workload.rs | 50 +++++++++++++++++++++++++++++++++--- 2 files changed, 47 insertions(+), 4 deletions(-) diff --git a/crates/types/src/types.rs b/crates/types/src/types.rs index 089ea9fa..7cd90bed 100644 --- a/crates/types/src/types.rs +++ b/crates/types/src/types.rs @@ -2,6 +2,7 @@ use serde::{Deserialize, Serialize}; /// Accuracy requirement that a query result must satisfy. #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] pub enum AccuracyTarget { /// Additive error bound ε: |estimate − true| ≤ ε · (domain size). Epsilon(f64), diff --git a/crates/types/src/workload.rs b/crates/types/src/workload.rs index 96ad2ca4..87dc4284 100644 --- a/crates/types/src/workload.rs +++ b/crates/types/src/workload.rs @@ -27,21 +27,26 @@ pub struct DurationMs(pub u64); /// SQL dialect variant — different dialects have different syntax and /// function sets that affect how the query string is parsed. #[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)] -#[serde(rename_all = "snake_case")] pub enum SqlDialect { + #[serde(rename = "datafusion_sql")] DataFusionSQL, + #[serde(rename = "clickhouse_sql")] ClickhouseSQL, + #[serde(rename = "elastic_sql")] ElasticSQL, } /// Source language of every query in the workload. /// All queries in a single `QueryWorkload` share the same language. #[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)] -#[serde(rename_all = "snake_case")] pub enum QueryLanguage { + #[serde(rename = "promql")] PromQL, + #[serde(rename = "sql")] SQL(SqlDialect), + #[serde(rename = "datafusion")] DataFusion, + #[serde(rename = "elastic_dsl")] ElasticDSL, } @@ -111,7 +116,7 @@ impl Default for QueryRequirements { /// let exploration = Predictability::AdHoc; /// ``` #[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] -#[serde(rename_all = "snake_case")] +#[serde(rename_all = "snake_case", deny_unknown_fields)] pub enum Predictability { AdHoc, Predictable { @@ -278,7 +283,7 @@ pub enum RepeatedDemand { /// ); /// ``` #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -#[serde(rename_all = "snake_case")] +#[serde(rename_all = "snake_case", deny_unknown_fields)] pub enum QueryRecurrence { OneTime { invocations: u64, @@ -807,6 +812,43 @@ mod tests { assert!(serde_json::from_str::(&with_unknown).is_err()); } + #[test] + fn query_language_wire_names_do_not_split_acronyms() { + let cases = [ + (QueryLanguage::PromQL, serde_json::json!("promql")), + ( + QueryLanguage::SQL(SqlDialect::DataFusionSQL), + serde_json::json!({"sql": "datafusion_sql"}), + ), + ( + QueryLanguage::SQL(SqlDialect::ClickhouseSQL), + serde_json::json!({"sql": "clickhouse_sql"}), + ), + (QueryLanguage::ElasticDSL, serde_json::json!("elastic_dsl")), + ]; + for (language, expected) in cases { + assert_eq!(serde_json::to_value(language).unwrap(), expected); + } + } + + #[test] + fn struct_enum_variants_reject_unknown_fields() { + assert!( + serde_json::from_value::(serde_json::json!({ + "one_time": {"invocations": 1, "execute_at": null, "typo": true} + })) + .is_err() + ); + assert!(serde_json::from_value::(serde_json::json!({ + "predictable": {"known_at": null, "typo": true} + })) + .is_err()); + assert!(serde_json::from_value::(serde_json::json!({ + "EpsilonDelta": {"epsilon": 0.1, "delta": 0.01, "typo": true} + })) + .is_err()); + } + #[test] fn estimated_demand_validates_window_rate_and_confidence() { let mut workload = base_workload();