From 766d96b48e0eccddb6594b304f35d655adc16f8f Mon Sep 17 00:00:00 2001 From: zz_y Date: Tue, 21 Jul 2026 09:55:56 -0600 Subject: [PATCH] feat(asap_types): retire WindowType in favor of asap_ir::WindowKind WindowType (Tumbling/Sliding) was a local duplicate of the same concept ASAPController's IR already models as WindowKind (Tumbling/Sliding/Session), just missing the Session variant and the serving-time traits (Copy/Default/Hash/Display/FromStr, snake_case serde) this workspace's real call sites need. ASAPController PR #143 added those upstream; this bumps the pinned asap-ir/asap-sketch/ asap-plan rev to the merge commit and re-exports WindowKind from asap_types::enums in WindowType's place. First-ever asap-ir dependency for asap_types (and transitively data_plane) -- previously only control_plane depended on ASAPController crates. WindowType lived in the shared crate and had ~35 real data_plane call sites, so unifying it for real (not just adding a Session variant locally) meant crossing that boundary. Pure rename at every call site -- same Tumbling default, same lowercase Display/FromStr round-trip, same wire format. No logic changed. Co-Authored-By: Claude Sonnet 5 --- Cargo.lock | 7 +-- control_plane/Cargo.toml | 16 +++---- control_plane/src/asap_tier_analysis.rs | 2 +- crates/asap_types/Cargo.toml | 6 +++ crates/asap_types/src/aggregation_config.rs | 10 ++--- crates/asap_types/src/enums.rs | 43 ++++++------------- crates/asap_types/src/policy_fingerprint.rs | 4 +- crates/asap_types/src/policy_registry.rs | 4 +- data_plane/benches/sketch_db.rs | 4 +- data_plane/src/drivers/ingest/otel.rs | 4 +- data_plane/src/drivers/query/servers/http.rs | 4 +- .../precompute_engine/accumulator_factory.rs | 12 +++--- .../src/precompute_engine/ingest_handler.rs | 4 +- .../src/precompute_engine/output_sink.rs | 4 +- data_plane/src/precompute_engine/worker.rs | 8 ++-- .../query_engines/asap_query_engine/engine.rs | 4 +- .../src/storage_engines/sketch_db/accuracy.rs | 4 +- .../sketch_db/backfill/processor.rs | 4 +- .../sketch_db/backfill/service.rs | 4 +- .../sketch_db/backfill/window_builder.rs | 4 +- .../sketch_db/lifecycle/eviction.rs | 4 +- .../sketch_db/lifecycle/reconcile.rs | 4 +- data_plane/src/storage_engines/types/enums.rs | 2 +- .../types/hot_reload_config.rs | 4 +- .../accuracy_empirical_validation_tests.rs | 4 +- .../tests/test_utilities/engine_factories.rs | 18 ++++---- .../tests/e2e_modified_otlp_sketch_path.rs | 12 +++--- 27 files changed, 94 insertions(+), 106 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index aef78304..82a26411 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -343,7 +343,7 @@ dependencies = [ [[package]] name = "asap-ir" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPController?rev=150ef7d0786d24286b578dfae9dbbcefc8fbac3e#150ef7d0786d24286b578dfae9dbbcefc8fbac3e" +source = "git+https://github.com/ProjectASAP/ASAPController?rev=7fcaf914d87e71407c3a6d7ccac613b867f9c11b#7fcaf914d87e71407c3a6d7ccac613b867f9c11b" dependencies = [ "serde", "serde_json", @@ -353,7 +353,7 @@ dependencies = [ [[package]] name = "asap-plan" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPController?rev=150ef7d0786d24286b578dfae9dbbcefc8fbac3e#150ef7d0786d24286b578dfae9dbbcefc8fbac3e" +source = "git+https://github.com/ProjectASAP/ASAPController?rev=7fcaf914d87e71407c3a6d7ccac613b867f9c11b#7fcaf914d87e71407c3a6d7ccac613b867f9c11b" dependencies = [ "asap-ir", "asap-sketch", @@ -373,7 +373,7 @@ dependencies = [ [[package]] name = "asap-sketch" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPController?rev=150ef7d0786d24286b578dfae9dbbcefc8fbac3e#150ef7d0786d24286b578dfae9dbbcefc8fbac3e" +source = "git+https://github.com/ProjectASAP/ASAPController?rev=7fcaf914d87e71407c3a6d7ccac613b867f9c11b#7fcaf914d87e71407c3a6d7ccac613b867f9c11b" dependencies = [ "asap-ir", ] @@ -409,6 +409,7 @@ name = "asap_types" version = "0.1.0" dependencies = [ "anyhow", + "asap-ir", "clap 4.6.1", "serde", "serde_json", diff --git a/control_plane/Cargo.toml b/control_plane/Cargo.toml index 4a72e880..354dd32d 100644 --- a/control_plane/Cargo.toml +++ b/control_plane/Cargo.toml @@ -39,14 +39,14 @@ asap_types.workspace = true # tagged releases yet. Re-pin as ASAPController's IR evolves; move to a tag # once one exists. # -# Bumped to 150ef7d (merge of PR #142, "derive PartialOrd/Ord for -# SummaryKind") for the sketch-identity unification work (see -# scratchpad/artifacts/enum-unification-plan.md) -- SummaryKind needs Ord -# for the BTreeSet deterministic-emission-order contract in -# control_plane::physical::colored_dag::emitter::EdgeStageConfig::metric_to_family. -asap-ir = { git = "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/ProjectASAP/ASAPController", rev = "150ef7d0786d24286b578dfae9dbbcefc8fbac3e" } -asap-sketch = { git = "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/ProjectASAP/ASAPController", rev = "150ef7d0786d24286b578dfae9dbbcefc8fbac3e" } -asap-plan = { git = "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/ProjectASAP/ASAPController", rev = "150ef7d0786d24286b578dfae9dbbcefc8fbac3e" } +# Bumped to 7fcaf91 (merge of PR #143, "add serving-time traits to +# WindowKind") for the WindowType -> WindowKind unification (see +# scratchpad/artifacts/enum-unification-plan.md) -- WindowKind needs +# Copy/Default/Hash/Display/FromStr + snake_case serde for +# asap_types::enums::WindowKind to replace the backend's own WindowType. +asap-ir = { git = "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/ProjectASAP/ASAPController", rev = "7fcaf914d87e71407c3a6d7ccac613b867f9c11b" } +asap-sketch = { git = "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/ProjectASAP/ASAPController", rev = "7fcaf914d87e71407c3a6d7ccac613b867f9c11b" } +asap-plan = { git = "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/ProjectASAP/ASAPController", rev = "7fcaf914d87e71407c3a6d7ccac613b867f9c11b" } [dev-dependencies] tokio = { version = "1", features = ["full", "test-util"] } diff --git a/control_plane/src/asap_tier_analysis.rs b/control_plane/src/asap_tier_analysis.rs index 4ef8dd74..1b85bc67 100644 --- a/control_plane/src/asap_tier_analysis.rs +++ b/control_plane/src/asap_tier_analysis.rs @@ -1442,7 +1442,7 @@ mod tests { String::new(), window_size, window_size, - asap_types::enums::WindowType::Tumbling, + asap_types::enums::WindowKind::Tumbling, spatial_filter.to_string(), metric.to_string(), None, diff --git a/crates/asap_types/Cargo.toml b/crates/asap_types/Cargo.toml index 31731c99..51550581 100644 --- a/crates/asap_types/Cargo.toml +++ b/crates/asap_types/Cargo.toml @@ -11,3 +11,9 @@ serde_yaml.workspace = true anyhow.workspace = true clap.workspace = true xxhash-rust = { version = "0.8", features = ["xxh64"] } + +# First cross-crate ASAPController dependency for data_plane (transitively, +# via this crate): WindowType -> asap_ir::intent_algebra::query_expr::WindowKind +# unification (scratchpad/artifacts/enum-unification-plan.md). Pin matches +# control_plane's -- see control_plane/Cargo.toml's comment for the rationale. +asap-ir = { git = "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/ProjectASAP/ASAPController", rev = "7fcaf914d87e71407c3a6d7ccac613b867f9c11b" } diff --git a/crates/asap_types/src/aggregation_config.rs b/crates/asap_types/src/aggregation_config.rs index a6c68dcc..e8d39a55 100644 --- a/crates/asap_types/src/aggregation_config.rs +++ b/crates/asap_types/src/aggregation_config.rs @@ -3,7 +3,7 @@ use serde_json::Value; use serde_yaml; use std::collections::HashMap; -use crate::enums::{QueryLanguage, WindowType}; +use crate::enums::{QueryLanguage, WindowKind}; use crate::policy_fingerprint::PolicyFingerprint; use crate::traits::SerializableToSink; use crate::utils::normalize_spatial_filter; @@ -33,7 +33,7 @@ pub struct AggregationConfig { pub window_size: u64, // Window size in seconds (e.g., 900s for 15m) pub slide_interval: u64, // Slide/hop interval in seconds (e.g., 30s) - pub window_type: WindowType, // Tumbling or Sliding + pub window_type: WindowKind, // Tumbling or Sliding pub spatial_filter: String, pub spatial_filter_normalized: String, @@ -86,7 +86,7 @@ impl AggregationConfig { original_yaml: String, window_size: u64, slide_interval: u64, - window_type: WindowType, + window_type: WindowKind, spatial_filter: String, metric: String, num_aggregates_to_retain: Option, @@ -176,7 +176,7 @@ impl AggregationConfig { .get("windowType") .and_then(|v| v.as_str()) .unwrap_or("tumbling") - .parse::() + .parse::() .unwrap_or_default(); let slide_interval = data @@ -296,7 +296,7 @@ impl AggregationConfig { .get("windowType") .and_then(|v| v.as_str()) .unwrap_or("tumbling") - .parse::() + .parse::() .unwrap_or_default(); let slide_interval = aggregation_data diff --git a/crates/asap_types/src/enums.rs b/crates/asap_types/src/enums.rs index e50bf57e..1fc7d063 100644 --- a/crates/asap_types/src/enums.rs +++ b/crates/asap_types/src/enums.rs @@ -116,34 +116,15 @@ impl FromStr for CleanupPolicy { } } -/// Window type for streaming aggregations. -#[derive( - Clone, Debug, Copy, Default, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize, -)] -#[serde(rename_all = "snake_case")] -pub enum WindowType { - #[default] - Tumbling, - Sliding, -} - -impl fmt::Display for WindowType { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - match self { - WindowType::Tumbling => write!(f, "tumbling"), - WindowType::Sliding => write!(f, "sliding"), - } - } -} - -impl FromStr for WindowType { - type Err = String; - - fn from_str(s: &str) -> Result { - match s.to_lowercase().as_str() { - "tumbling" => Ok(WindowType::Tumbling), - "sliding" => Ok(WindowType::Sliding), - _ => Err(format!("Unknown window type: '{s}'")), - } - } -} +/// Window lifecycle/flush semantics for streaming aggregations. +/// +/// Formerly a local `WindowType` (`Tumbling`/`Sliding`) enum. Retired in +/// favor of `asap_ir::intent_algebra::query_expr::WindowKind` directly — +/// same concept, plus a `Session` variant this workspace didn't have. +/// `Copy`/`Default`/`Hash`/`Display`/`FromStr` and +/// `#[serde(rename_all = "snake_case")]` were added upstream +/// (ASAPController PR #143) specifically so this re-export could replace +/// the old local type without touching any call site's behavior: same +/// `Tumbling` default, same lowercase `Display`/`FromStr` round-trip, same +/// wire format. +pub use asap_ir::intent_algebra::query_expr::WindowKind; diff --git a/crates/asap_types/src/policy_fingerprint.rs b/crates/asap_types/src/policy_fingerprint.rs index 7a5c285f..fcaea645 100644 --- a/crates/asap_types/src/policy_fingerprint.rs +++ b/crates/asap_types/src/policy_fingerprint.rs @@ -172,7 +172,7 @@ impl std::fmt::Display for PolicyFingerprint { #[cfg(test)] mod tests { use super::*; - use crate::enums::WindowType; + use crate::enums::WindowKind; use crate::AggregationType; use crate::KeyByLabelNames; use std::collections::HashMap; @@ -195,7 +195,7 @@ mod tests { String::new(), window_size, window_size, - WindowType::Tumbling, + WindowKind::Tumbling, spatial_filter.to_string(), metric.to_string(), None, diff --git a/crates/asap_types/src/policy_registry.rs b/crates/asap_types/src/policy_registry.rs index 99cd399d..a10f0548 100644 --- a/crates/asap_types/src/policy_registry.rs +++ b/crates/asap_types/src/policy_registry.rs @@ -115,7 +115,7 @@ impl PolicyRegistry { #[cfg(test)] mod tests { use super::*; - use crate::enums::WindowType; + use crate::enums::WindowKind; use crate::AggregationType; use crate::KeyByLabelNames; use std::collections::HashMap as StdHashMap; @@ -134,7 +134,7 @@ mod tests { String::new(), 60, 60, - WindowType::Tumbling, + WindowKind::Tumbling, String::new(), metric.to_string(), None, diff --git a/data_plane/benches/sketch_db.rs b/data_plane/benches/sketch_db.rs index 1ef08fcd..845e574f 100644 --- a/data_plane/benches/sketch_db.rs +++ b/data_plane/benches/sketch_db.rs @@ -388,7 +388,7 @@ fn bench_query_precomputes_by_agg(c: &mut Criterion) { /// case, where the per-batch reconcile is pure scan overhead. fn matching_streaming_config(metric: &str) -> data_plane::storage_engines::types::StreamingConfig { use asap_types::aggregation_config::AggregationConfig; - use asap_types::enums::WindowType; + use asap_types::enums::WindowKind; use asap_types::AggregationType as AT; use asap_types::KeyByLabelNames; use std::collections::HashMap; @@ -403,7 +403,7 @@ fn matching_streaming_config(metric: &str) -> data_plane::storage_engines::types String::new(), 60, 60, - WindowType::Tumbling, + WindowKind::Tumbling, String::new(), metric.to_string(), None, diff --git a/data_plane/src/drivers/ingest/otel.rs b/data_plane/src/drivers/ingest/otel.rs index aca607c4..305556fc 100644 --- a/data_plane/src/drivers/ingest/otel.rs +++ b/data_plane/src/drivers/ingest/otel.rs @@ -4003,7 +4003,7 @@ mod sid_bucketing_tests { Metric as PbMetric, NumberDataPoint, ResourceMetrics, ScopeMetrics, }; use asap_types::aggregation_config::AggregationConfig; - use asap_types::enums::WindowType; + use asap_types::enums::WindowKind; use asap_types::AggregationType; use asap_types::KeyByLabelNames; use std::collections::HashMap; @@ -4030,7 +4030,7 @@ mod sid_bucketing_tests { String::new(), 10, 10, - WindowType::Tumbling, + WindowKind::Tumbling, String::new(), metric.to_string(), None, diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index 13595ae2..eba747b2 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -3044,7 +3044,7 @@ aggregations: active_agg_ids: &[u64], ) -> (u16, std::collections::HashMap) { use asap_types::aggregation_config::AggregationConfig; - use asap_types::enums::WindowType; + use asap_types::enums::WindowKind; use asap_types::AggregationType; use asap_types::KeyByLabelNames; use std::collections::HashMap; @@ -3075,7 +3075,7 @@ aggregations: original_yaml: String::new(), window_size: 1, slide_interval: 1, - window_type: WindowType::Tumbling, + window_type: WindowKind::Tumbling, spatial_filter: String::new(), spatial_filter_normalized: String::new(), metric: metric.clone(), diff --git a/data_plane/src/precompute_engine/accumulator_factory.rs b/data_plane/src/precompute_engine/accumulator_factory.rs index 0766a4a9..aacd8e76 100644 --- a/data_plane/src/precompute_engine/accumulator_factory.rs +++ b/data_plane/src/precompute_engine/accumulator_factory.rs @@ -948,7 +948,7 @@ pub fn create_accumulator_updater(config: &AggregationConfig) -> Box) -> Aggregation String::new(), 60, 60, - WindowType::Tumbling, + WindowKind::Tumbling, String::new(), "m".to_string(), None, diff --git a/data_plane/src/tests/test_utilities/engine_factories.rs b/data_plane/src/tests/test_utilities/engine_factories.rs index c71c326e..c7d42eb0 100644 --- a/data_plane/src/tests/test_utilities/engine_factories.rs +++ b/data_plane/src/tests/test_utilities/engine_factories.rs @@ -10,7 +10,7 @@ use crate::query_engines::asap_query_engine::engine::ASAPQueryEngine; use crate::query_engines::query_result::InstantVectorElement; use crate::storage_engines::types::{ AggregationConfig, AggregationType, KeyByLabelValues, PrecomputedOutput, QueryLanguage, - StreamingConfig, WindowType, + StreamingConfig, WindowKind, }; use crate::AggregateCore; use asap_types::KeyByLabelNames; @@ -99,7 +99,7 @@ pub fn create_engine_single_pop_with_aggregated( original_yaml: String::new(), window_size: 1, slide_interval: 1, - window_type: WindowType::Tumbling, + window_type: WindowKind::Tumbling, spatial_filter: String::new(), spatial_filter_normalized: String::new(), metric: metric.to_string(), @@ -191,7 +191,7 @@ pub fn create_engine_dual_input( original_yaml: String::new(), window_size: 1, slide_interval: 1, - window_type: WindowType::Tumbling, + window_type: WindowKind::Tumbling, spatial_filter: String::new(), spatial_filter_normalized: String::new(), metric: metric.to_string(), @@ -213,7 +213,7 @@ pub fn create_engine_dual_input( original_yaml: String::new(), window_size: 1, slide_interval: 1, - window_type: WindowType::Tumbling, + window_type: WindowKind::Tumbling, spatial_filter: String::new(), spatial_filter_normalized: String::new(), metric: metric.to_string(), @@ -300,7 +300,7 @@ pub fn create_engine_two_metrics( original_yaml: String::new(), window_size: 1, slide_interval: 1, - window_type: WindowType::Tumbling, + window_type: WindowKind::Tumbling, spatial_filter: String::new(), spatial_filter_normalized: String::new(), metric: metric_a.to_string(), @@ -321,7 +321,7 @@ pub fn create_engine_two_metrics( original_yaml: String::new(), window_size: 1, slide_interval: 1, - window_type: WindowType::Tumbling, + window_type: WindowKind::Tumbling, spatial_filter: String::new(), spatial_filter_normalized: String::new(), metric: metric_b.to_string(), @@ -418,7 +418,7 @@ pub fn create_engine_three_metrics( original_yaml: String::new(), window_size: 1, slide_interval: 1, - window_type: WindowType::Tumbling, + window_type: WindowKind::Tumbling, spatial_filter: String::new(), spatial_filter_normalized: String::new(), metric: metric.to_string(), @@ -492,7 +492,7 @@ pub fn create_engine_multi_timestamp( original_yaml: String::new(), window_size: 1, slide_interval: 1, - window_type: WindowType::Tumbling, + window_type: WindowKind::Tumbling, spatial_filter: String::new(), spatial_filter_normalized: String::new(), metric: metric.to_string(), @@ -542,7 +542,7 @@ pub fn create_engine_multi_timestamp_with_window( data: Vec<(u64, Option>, Box)>, promql_query: &str, window_size: u64, - window_type: WindowType, + window_type: WindowKind, ) -> ASAPQueryEngine { let grouping_label_strings: Vec = grouping_labels.iter().map(|s| s.to_string()).collect(); diff --git a/data_plane/tests/e2e_modified_otlp_sketch_path.rs b/data_plane/tests/e2e_modified_otlp_sketch_path.rs index 25089f2b..fa2ed1a4 100644 --- a/data_plane/tests/e2e_modified_otlp_sketch_path.rs +++ b/data_plane/tests/e2e_modified_otlp_sketch_path.rs @@ -38,7 +38,7 @@ use asap_sketchlib::proto::sketchlib::{ }; use asap_sketchlib::MessagePackCodec; use asap_types::aggregation_config::AggregationConfig; -use asap_types::enums::WindowType; +use asap_types::enums::WindowKind; use asap_types::AggregationType; use prost::Message; use std::collections::HashMap; @@ -78,7 +78,7 @@ fn make_count_min_agg_config( String::new(), window_secs, 0, - WindowType::Tumbling, + WindowKind::Tumbling, metric.to_string(), metric.to_string(), None, @@ -344,7 +344,7 @@ fn make_count_sketch_agg_config( String::new(), window_secs, 0, - WindowType::Tumbling, + WindowKind::Tumbling, metric.to_string(), metric.to_string(), None, @@ -555,7 +555,7 @@ fn make_kll_agg_config( String::new(), window_secs, 0, - WindowType::Tumbling, + WindowKind::Tumbling, metric.to_string(), metric.to_string(), None, @@ -735,7 +735,7 @@ fn make_dd_sketch_agg_config( String::new(), window_secs, 0, - WindowType::Tumbling, + WindowKind::Tumbling, metric.to_string(), metric.to_string(), None, @@ -907,7 +907,7 @@ fn make_hll_agg_config( String::new(), window_secs, 0, - WindowType::Tumbling, + WindowKind::Tumbling, metric.to_string(), metric.to_string(), None,