diff --git a/crates/asap_types/src/aggregation_config.rs b/crates/asap_types/src/aggregation_config.rs index 24d804d2..f256d74f 100644 --- a/crates/asap_types/src/aggregation_config.rs +++ b/crates/asap_types/src/aggregation_config.rs @@ -321,7 +321,7 @@ impl PrecomputeMaterialization { /// `PolicyFingerprint::as_u64()` — the u64-form handle used by the /// policy-fingerprint-keyed call sites (e.g. `InstalledPrecomputePlan`'s - /// `HashMap` keys). **Always** equal to + /// `HashMap` keys). **Always** equal to /// `self.policy_fingerprint().as_u64()`. The value is content- /// addressed identity, NOT a controller-allocated counter id. pub fn policy_fp_u64(&self) -> u64 { diff --git a/crates/asap_types/src/policy_registry.rs b/crates/asap_types/src/policy_registry.rs index 16f007bb..90138b13 100644 --- a/crates/asap_types/src/policy_registry.rs +++ b/crates/asap_types/src/policy_registry.rs @@ -22,7 +22,7 @@ //! //! Two `PrecomputeMaterialization`s that produce the same `PolicyFingerprint` //! ARE the same policy. The registry treats this as a *deduplication* -//! invariant — if two distinct entries in the source `materializations_by_policy_fingerprint` +//! invariant — if two distinct entries in the source `materializations_by_output` //! map produce the same fingerprint, the later one wins (last-write //! semantics). In practice the source should never contain duplicates; //! if it does, that's a control-plane bug worth surfacing in telemetry diff --git a/crates/asap_types/src/precompute_plan.rs b/crates/asap_types/src/precompute_plan.rs index f360ab4b..05227b67 100644 --- a/crates/asap_types/src/precompute_plan.rs +++ b/crates/asap_types/src/precompute_plan.rs @@ -422,13 +422,21 @@ impl PrecomputePlan { pub fn runtime_materializations( &self, - ) -> Result, PrecomputePlanError> { + ) -> Result< + HashMap, + PrecomputePlanError, + > { self.validate()?; Ok(self .materializations .iter() .cloned() - .map(|materialization| (materialization.policy_fp_u64(), materialization)) + .map(|materialization| { + ( + crate::sds::StoredOutputId(materialization.policy_fp_u64()), + materialization, + ) + }) .collect()) } diff --git a/crates/asap_types/src/sds.rs b/crates/asap_types/src/sds.rs index 2e665f9c..a6870ea2 100644 --- a/crates/asap_types/src/sds.rs +++ b/crates/asap_types/src/sds.rs @@ -43,6 +43,11 @@ impl StoredOutputId { self.0 } } +impl From for StoredOutputId { + fn from(value: u64) -> Self { + Self(value) + } +} impl From for StoredOutputId { fn from(value: crate::PolicyFingerprint) -> Self { Self(value.0) diff --git a/data_plane/src/drivers/ingest/otel.rs b/data_plane/src/drivers/ingest/otel.rs index 7185b479..81925b1a 100644 --- a/data_plane/src/drivers/ingest/otel.rs +++ b/data_plane/src/drivers/ingest/otel.rs @@ -1405,7 +1405,7 @@ async fn route_modified_otlp_sketches_to_precompute( // below so the query engine can answer per-item estimate(key) // (the CMS/CountSketch FrequencyEstimate gate consults it). let item_label_for_sid: Option = { - snap.get_aggregation_config(policy_fp.as_u64()) + snap.get_aggregation_config(policy_fp.into()) .or_else(|| { snap.materializations() .values() @@ -4621,7 +4621,7 @@ mod sid_bucketing_tests { let policy_fp = asap_types::PolicyFingerprint(cfg.policy_fp_u64()); let mut configs = HashMap::new(); configs.insert(cfg.policy_fp_u64(), cfg.clone()); - let streaming = InstalledPrecomputePlan::new(configs); + let streaming = InstalledPrecomputePlan::from_raw_ids(configs); let hot_reload = InstalledPrecomputePlanHandle::new(streaming); let resolver = Arc::new(SeriesIdResolver::new()); diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index 3816c0db..258c66cd 100644 --- a/data_plane/src/drivers/ingest/prometheus_remote_write.rs +++ b/data_plane/src/drivers/ingest/prometheus_remote_write.rs @@ -435,7 +435,7 @@ impl PrometheusRemoteWriteReceiver { continue; }; let config = snapshot - .get_aggregation_config(policy_fp.as_u64()) + .get_aggregation_config((*policy_fp).into()) .ok_or(RemoteWriteError::InactivePhysicalPlan)?; let manager = crate::precompute_engine::window_manager::WindowManager::with_layout( config.window_size, @@ -987,7 +987,7 @@ mod tests { capability_snapshot_id: "test".into(), }; let configs = streaming - .materializations_by_policy_fingerprint + .materializations_by_output .values() .cloned() .collect::>(); @@ -1021,7 +1021,7 @@ mod tests { producers: Vec::new(), executable_dags: Default::default(), materializations: streaming - .materializations_by_policy_fingerprint + .materializations_by_output .values() .cloned() .collect(), @@ -1094,7 +1094,8 @@ mod tests { value_source_column: None, }; let policy_fp = aggregation.policy_fp_u64(); - let streaming = InstalledPrecomputePlan::new(HashMap::from([(policy_fp, aggregation)])); + let streaming = + InstalledPrecomputePlan::from_raw_ids(HashMap::from([(policy_fp, aggregation)])); let (sender, receiver) = mpsc::channel(8); let ingest = Arc::new(IngestState { router: SeriesRouter::new(vec![sender]), @@ -1132,7 +1133,7 @@ mod tests { let mut config = snapshot.precompute_plan.materializations[0].clone(); config.population_key_encoding = asap_types::PopulationKeyEncoding::CanonicalLabelsV1; config.partitioning = Some(asap_types::sds::PopulationPartitioning::Grouped); - let hot = physical_config(InstalledPrecomputePlan::new(HashMap::from([( + let hot = physical_config(InstalledPrecomputePlan::from_raw_ids(HashMap::from([( config.policy_fp_u64(), config.clone(), )]))); @@ -1248,7 +1249,7 @@ mod tests { assert_ne!(kll_fp, pooled_kll_fp); let cms_fp = cms.policy_fingerprint(); let counter_fp = counter.policy_fingerprint(); - let streaming = InstalledPrecomputePlan::new(HashMap::from([ + let streaming = InstalledPrecomputePlan::from_raw_ids(HashMap::from([ (cms_fp.0, cms), (counter_fp.0, counter), (kll_fp.0, kll), @@ -1652,18 +1653,18 @@ mod tests { let policy = *ingest .hot_reload_config .snapshot() - .materializations_by_policy_fingerprint + .materializations_by_output .keys() .next() .unwrap(); let binding = asap_types::query_plan::MaterializationBinding { full_window_slide_ms: None, - materialization: asap_types::PolicyFingerprint(policy).into(), + materialization: policy, stored_output_reference: ingest .summary_store .summary_catalog_snapshot() .unwrap() - .output_reference(asap_types::PolicyFingerprint(policy).into()) + .output_reference(policy) .unwrap(), output_grouping: asap_types::query_plan::PhysicalGrouping::Reduce(vec!["job".into()]), item_labels: vec![], diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index 1199e886..5480a76f 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -3274,7 +3274,7 @@ mod tests { marker_to_fp.insert(*marker, fp); agg_map.insert(fp, cfg); } - let installed_precompute_plan = Arc::new(InstalledPrecomputePlan::new(agg_map)); + let installed_precompute_plan = Arc::new(InstalledPrecomputePlan::from_raw_ids(agg_map)); let hot_reload = InstalledPrecomputePlanHandle::from_arc(installed_precompute_plan.clone()); let query_engine = Arc::new(ASAPQueryEngine::new(15000)); let summary_store = Arc::new(crate::storage_engines::sketch_db::index::SketchStore::new()); @@ -5418,10 +5418,10 @@ pub fn validate_and_build_runtime_plan( ) .map_err(|error| format!("DAG execution installation failed: {error}"))?; let typed_fps: BTreeSet<_> = installed_precompute_plan - .materializations_by_policy_fingerprint + .materializations_by_output .keys() .copied() - .map(asap_types::PolicyFingerprint) + .map(asap_types::sds::StoredOutputId::fingerprint) .collect(); request .query_plan @@ -6245,7 +6245,8 @@ async fn handle_post_backfill_job( return (StatusCode::SERVICE_UNAVAILABLE, axum::Json(body)).into_response(); }; let snapshot = handle.snapshot(); - let agg_cfg = match snapshot.get_aggregation_config(req.agg_id) { + let agg_cfg = match snapshot.get_aggregation_config(asap_types::sds::StoredOutputId(req.agg_id)) + { Some(c) => c.clone(), None => { let body = serde_json::json!({ diff --git a/data_plane/src/precompute_engine/ingest_handler.rs b/data_plane/src/precompute_engine/ingest_handler.rs index cf033c14..400797a7 100644 --- a/data_plane/src/precompute_engine/ingest_handler.rs +++ b/data_plane/src/precompute_engine/ingest_handler.rs @@ -336,7 +336,7 @@ mod tests { let mut map = std::collections::HashMap::new(); map.insert(agg_id, make_config(agg_id, metric)); - let streaming = InstalledPrecomputePlan::new(map); + let streaming = InstalledPrecomputePlan::from_raw_ids(map); let hot_reload = crate::storage_engines::types::InstalledPrecomputePlanHandle::new(streaming.clone()); diff --git a/data_plane/src/precompute_engine/output_sink.rs b/data_plane/src/precompute_engine/output_sink.rs index 42f2562c..eab97b32 100644 --- a/data_plane/src/precompute_engine/output_sink.rs +++ b/data_plane/src/precompute_engine/output_sink.rs @@ -415,7 +415,7 @@ mod tests { let agg_id = cfg.policy_fp_u64(); let mut configs = HashMap::new(); configs.insert(agg_id, cfg); - let streaming = InstalledPrecomputePlan::new(configs); + let streaming = InstalledPrecomputePlan::from_raw_ids(configs); let hot_reload = InstalledPrecomputePlanHandle::new(streaming.clone()); let summary_store = Arc::new(SketchStore::new()); @@ -484,10 +484,9 @@ mod tests { Arc::new(SeriesIdResolver::open(temporary.path().join("resolver.wal")).unwrap()); let sink = SketchStoreSink::new( store.clone(), - InstalledPrecomputePlanHandle::new(InstalledPrecomputePlan::new(HashMap::from([( - fingerprint.0, - cfg, - )]))), + InstalledPrecomputePlanHandle::new(InstalledPrecomputePlan::from_raw_ids( + HashMap::from([(fingerprint.0, cfg)]), + )), resolver, ); let original_generation = Arc::new(catalog.reference().unwrap()); @@ -566,7 +565,7 @@ mod tests { cfg.parameters .insert("alpha".into(), serde_json::json!(0.01)); let policy_fp = cfg.policy_fp_u64(); - let hot_reload = InstalledPrecomputePlanHandle::new(InstalledPrecomputePlan::new( + let hot_reload = InstalledPrecomputePlanHandle::new(InstalledPrecomputePlan::from_raw_ids( HashMap::from([(policy_fp, cfg)]), )); let summary_store = Arc::new(SketchStore::new()); diff --git a/data_plane/src/precompute_engine/series_router.rs b/data_plane/src/precompute_engine/series_router.rs index 6a29e4cf..0334f5fb 100644 --- a/data_plane/src/precompute_engine/series_router.rs +++ b/data_plane/src/precompute_engine/series_router.rs @@ -40,7 +40,7 @@ pub enum WorkerMessage { sid: u64, /// Source `PrecomputeMaterialization` fingerprint. Worker looks up its /// `PrecomputeMaterialization` (window size, sketch kind/config, late - /// data policy, etc.) via `snap.get_aggregation_config(policy_fp.as_u64())`. + /// data policy, etc.) via `snap.get_aggregation_config(policy_fp.into())`. policy_fp: PolicyFingerprint, /// Grouping label values joined by semicolons (e.g. "constant"). /// Empty string if the aggregation has no grouping labels. Used diff --git a/data_plane/src/precompute_engine/worker.rs b/data_plane/src/precompute_engine/worker.rs index d957c0bf..6e523017 100644 --- a/data_plane/src/precompute_engine/worker.rs +++ b/data_plane/src/precompute_engine/worker.rs @@ -513,7 +513,7 @@ impl Worker { ) -> Result, String> { if !self.group_states.contains_key(&sid) { let snap = self.hot_reload.snapshot(); - let Some(cfg) = snap.get_aggregation_config(policy_fp.as_u64()) else { + let Some(cfg) = snap.get_aggregation_config(policy_fp.into()) else { return Ok(None); }; let program = snap.raw_programs.get(&policy_fp.as_u64()).cloned(); @@ -1081,7 +1081,7 @@ impl Worker { let snap = self.hot_reload.snapshot(); let before = self.group_states.len(); self.group_states.retain(|&sid, gs| { - if snap.contains(gs.policy_fp.as_u64()) { + if snap.contains(gs.policy_fp.into()) { return true; // policy still in config, keep } // Policy retired — keep only if there's residual data @@ -1976,6 +1976,7 @@ mod tests { use asap_physical_operators::summary_kernels::sum::SumAccumulator; use asap_sketchlib::KllSketch; use asap_types::enums::WindowKind; + use asap_types::sds::StoredOutputId; use asap_types::AggregationType; fn make_agg_config( @@ -2086,7 +2087,7 @@ mod tests { configs: HashMap, ) -> crate::storage_engines::types::InstalledPrecomputePlanHandle { crate::storage_engines::types::InstalledPrecomputePlanHandle::new( - crate::storage_engines::types::InstalledPrecomputePlan::new(configs), + crate::storage_engines::types::InstalledPrecomputePlan::from_raw_ids(configs), ) } diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs index e50ae3f0..c8162ffb 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs @@ -1046,10 +1046,9 @@ mod tests { name: "value".into(), }); let hot = crate::storage_engines::types::InstalledPrecomputePlanHandle::from_arc(Arc::new( - crate::storage_engines::types::InstalledPrecomputePlan::new(HashMap::from([( - cfg.policy_fp_u64(), - cfg.clone(), - )])), + crate::storage_engines::types::InstalledPrecomputePlan::from_raw_ids(HashMap::from([ + (cfg.policy_fp_u64(), cfg.clone()), + ])), )); let registry = Arc::new(crate::storage_engines::sketch_db::backfill::BackfillRegistry::new()); diff --git a/data_plane/src/storage_engines/sketch_db/backfill/processor.rs b/data_plane/src/storage_engines/sketch_db/backfill/processor.rs index ceb6e984..f72eaa48 100644 --- a/data_plane/src/storage_engines/sketch_db/backfill/processor.rs +++ b/data_plane/src/storage_engines/sketch_db/backfill/processor.rs @@ -212,7 +212,7 @@ impl WindowProcessor for BackfillWindowProcessor { ) -> Result<(), Box> { let snapshot = self.config.snapshot(); let config = snapshot - .get_aggregation_config(agg_id) + .get_aggregation_config(asap_types::sds::StoredOutputId(agg_id)) .cloned() .ok_or_else(|| format!("agg_id {agg_id} not in current InstalledPrecomputePlan"))?; let program = snapshot.raw_programs.get(&agg_id).cloned(); @@ -444,7 +444,7 @@ mod tests { fn streaming_config_with(config: PrecomputeMaterialization) -> Arc { let mut map = std::collections::HashMap::new(); map.insert(config.policy_fp_u64(), config); - Arc::new(InstalledPrecomputePlan::new(map)) + Arc::new(InstalledPrecomputePlan::from_raw_ids(map)) } #[tokio::test] diff --git a/data_plane/src/storage_engines/sketch_db/backfill/service.rs b/data_plane/src/storage_engines/sketch_db/backfill/service.rs index 35c9b135..8e563942 100644 --- a/data_plane/src/storage_engines/sketch_db/backfill/service.rs +++ b/data_plane/src/storage_engines/sketch_db/backfill/service.rs @@ -200,7 +200,9 @@ impl BackfillService { ); let snapshot = self.config_source.snapshot(); - let Some(materialization) = snapshot.get_aggregation_config(job.agg_id) else { + let Some(materialization) = + snapshot.get_aggregation_config(asap_types::sds::StoredOutputId(job.agg_id)) + else { self.registry .mark_failed(job.job_id, "backfill materialization is not installed"); continue; @@ -372,7 +374,7 @@ mod tests { fn streaming_with(cfg: PrecomputeMaterialization) -> Arc { let mut m = std::collections::HashMap::new(); m.insert(cfg.policy_fp_u64(), cfg); - Arc::new(InstalledPrecomputePlan::new(m)) + Arc::new(InstalledPrecomputePlan::from_raw_ids(m)) } async fn wait_for_status( diff --git a/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs b/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs index b3edf16f..4788aaca 100644 --- a/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs +++ b/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs @@ -244,7 +244,10 @@ mod tests { id_to_fp.insert(id, fp); map.insert(fp, cfg); } - (Arc::new(InstalledPrecomputePlan::new(map)), id_to_fp) + ( + Arc::new(InstalledPrecomputePlan::from_raw_ids(map)), + id_to_fp, + ) } fn write_one( @@ -263,7 +266,7 @@ mod tests { asap_types::PolicyFingerprint(agg_id), ); let agg_cfg = installed_precompute_plan - .get_aggregation_config(agg_id) + .get_aggregation_config(asap_types::sds::StoredOutputId(agg_id)) .unwrap(); // Test-scoped resolver — each call mints fresh. Production // shares one resolver across all sinks; tests don't need that diff --git a/data_plane/src/storage_engines/sketch_db/lifecycle/reconcile.rs b/data_plane/src/storage_engines/sketch_db/lifecycle/reconcile.rs index c5bd2d50..d6310b46 100644 --- a/data_plane/src/storage_engines/sketch_db/lifecycle/reconcile.rs +++ b/data_plane/src/storage_engines/sketch_db/lifecycle/reconcile.rs @@ -326,7 +326,7 @@ mod tests { for (i, c) in configs.into_iter().enumerate() { map.insert(i as u64 + 1, c); } - InstalledPrecomputePlan::new(map) + InstalledPrecomputePlan::from_raw_ids(map) } #[test] diff --git a/data_plane/src/storage_engines/types/hot_reload_config.rs b/data_plane/src/storage_engines/types/hot_reload_config.rs index 567852c8..0e37ac14 100644 --- a/data_plane/src/storage_engines/types/hot_reload_config.rs +++ b/data_plane/src/storage_engines/types/hot_reload_config.rs @@ -560,10 +560,7 @@ impl std::fmt::Debug for InstalledPrecomputePlanHandle { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { let snap = self.snapshot(); f.debug_struct("InstalledPrecomputePlanHandle") - .field( - "num_agg_configs", - &snap.materializations_by_policy_fingerprint.len(), - ) + .field("num_agg_configs", &snap.materializations_by_output.len()) .finish() } } @@ -711,7 +708,7 @@ mod tests { id_to_fp.insert(id, fp); map.insert(fp, cfg); } - (InstalledPrecomputePlan::new(map), id_to_fp) + (InstalledPrecomputePlan::from_raw_ids(map), id_to_fp) } #[test] @@ -719,10 +716,10 @@ mod tests { let (cfg, id_to_fp) = cfg_with_ids(&[1, 2, 3]); let hr = InstalledPrecomputePlanHandle::new(cfg); let snap = hr.snapshot(); - assert_eq!(snap.materializations_by_policy_fingerprint.len(), 3); + assert_eq!(snap.materializations_by_output.len(), 3); assert!(snap - .materializations_by_policy_fingerprint - .contains_key(&id_to_fp[&2])); + .materializations_by_output + .contains_key(&asap_types::sds::StoredOutputId(id_to_fp[&2]))); } #[test] @@ -732,19 +729,19 @@ mod tests { let hr = InstalledPrecomputePlanHandle::new(cfg1); let old = hr.swap(cfg2); // Old snapshot still reflects pre-swap contents. - assert_eq!(old.materializations_by_policy_fingerprint.len(), 2); + assert_eq!(old.materializations_by_output.len(), 2); assert!(old - .materializations_by_policy_fingerprint - .contains_key(&id_to_fp1[&1])); + .materializations_by_output + .contains_key(&asap_types::sds::StoredOutputId(id_to_fp1[&1]))); // New snapshot reflects post-swap contents. let new_snap = hr.snapshot(); - assert_eq!(new_snap.materializations_by_policy_fingerprint.len(), 3); + assert_eq!(new_snap.materializations_by_output.len(), 3); assert!(new_snap - .materializations_by_policy_fingerprint - .contains_key(&id_to_fp2[&5])); + .materializations_by_output + .contains_key(&asap_types::sds::StoredOutputId(id_to_fp2[&5]))); assert!(!new_snap - .materializations_by_policy_fingerprint - .contains_key(&id_to_fp1[&1])); + .materializations_by_output + .contains_key(&asap_types::sds::StoredOutputId(id_to_fp1[&1]))); } #[test] @@ -757,10 +754,10 @@ mod tests { // The clone sees the swap because both handles share the // same ArcSwap inside. let snap = hr_clone.snapshot(); - assert_eq!(snap.materializations_by_policy_fingerprint.len(), 2); + assert_eq!(snap.materializations_by_output.len(), 2); assert!(snap - .materializations_by_policy_fingerprint - .contains_key(&id_to_fp2[&3])); + .materializations_by_output + .contains_key(&asap_types::sds::StoredOutputId(id_to_fp2[&3]))); } #[test] @@ -781,7 +778,7 @@ mod tests { // Under race, the snapshot must be internally // consistent — either 2 entries (original) or 3 // (post-swap). Never a torn state. - let n = snap.materializations_by_policy_fingerprint.len(); + let n = snap.materializations_by_output.len(); assert!(n == 2 || n == 3, "torn snapshot: {n} entries"); } }); diff --git a/data_plane/src/storage_engines/types/installed_precompute_plan.rs b/data_plane/src/storage_engines/types/installed_precompute_plan.rs index 4ca67770..1ec25302 100644 --- a/data_plane/src/storage_engines/types/installed_precompute_plan.rs +++ b/data_plane/src/storage_engines/types/installed_precompute_plan.rs @@ -4,6 +4,7 @@ use std::collections::HashMap; use std::ops::Index; use super::storage_backend::StorageBackend; +use asap_types::sds::StoredOutputId; use asap_types::{PolicyRegistry, PrecomputeMaterialization}; #[derive(Debug, Clone)] @@ -12,17 +13,17 @@ pub struct InstalledPrecomputePlan { pub(crate) raw_programs: HashMap>, pub(crate) precompute_plan: Option, - pub(crate) materializations_by_policy_fingerprint: HashMap, + pub(crate) materializations_by_output: HashMap, pub(crate) storage_backend: StorageBackend, } impl InstalledPrecomputePlan { - fn derived_view(materializations: HashMap) -> Self { + fn derived_view(materializations: HashMap) -> Self { Self { partitioning: Default::default(), raw_programs: HashMap::new(), precompute_plan: None, - materializations_by_policy_fingerprint: materializations, + materializations_by_output: materializations, storage_backend: StorageBackend::default(), } } @@ -52,13 +53,26 @@ impl InstalledPrecomputePlan { // Isolated kernel/storage fixtures can omit a physical installation. This // constructor is absent from the production library and binary. #[cfg(test)] - pub fn new(materializations: HashMap) -> Self { + pub fn new(materializations: HashMap) -> Self { Self::derived_view(materializations) } + /// Fixtures hold raw output ids, so say that once here rather than + /// wrapping every literal. Production construction goes through + /// `from_precompute_plan`, which takes the ids from the plan itself. + #[cfg(test)] + pub fn from_raw_ids(materializations: HashMap) -> Self { + Self::new( + materializations + .into_iter() + .map(|(id, value)| (StoredOutputId(id), value)) + .collect(), + ) + } + #[cfg(test)] pub fn with_storage_backend( - materializations: HashMap, + materializations: HashMap, storage_backend: StorageBackend, ) -> Self { let mut view = Self::derived_view(materializations); @@ -79,7 +93,7 @@ impl InstalledPrecomputePlan { #[cfg(test)] let selected = selected.or_else(|| { let materializations = self - .materializations_by_policy_fingerprint + .materializations_by_output .values() .cloned() .collect::>(); @@ -105,33 +119,30 @@ impl InstalledPrecomputePlan { self.storage_backend } - pub fn get_aggregation_config(&self, fingerprint: u64) -> Option<&PrecomputeMaterialization> { - self.materializations_by_policy_fingerprint - .get(&fingerprint) + pub fn get_aggregation_config( + &self, + output: StoredOutputId, + ) -> Option<&PrecomputeMaterialization> { + self.materializations_by_output.get(&output) } - pub fn materializations(&self) -> &HashMap { - &self.materializations_by_policy_fingerprint + pub fn materializations(&self) -> &HashMap { + &self.materializations_by_output } - pub fn contains(&self, fingerprint: u64) -> bool { - self.materializations_by_policy_fingerprint - .contains_key(&fingerprint) + pub fn contains(&self, output: StoredOutputId) -> bool { + self.materializations_by_output.contains_key(&output) } pub fn policy_registry(&self) -> PolicyRegistry { - PolicyRegistry::from_configs( - self.materializations_by_policy_fingerprint - .values() - .cloned(), - ) + PolicyRegistry::from_configs(self.materializations_by_output.values().cloned()) } } -impl Index for InstalledPrecomputePlan { +impl Index for InstalledPrecomputePlan { type Output = PrecomputeMaterialization; - fn index(&self, fingerprint: u64) -> &Self::Output { - &self.materializations_by_policy_fingerprint[&fingerprint] + fn index(&self, output: StoredOutputId) -> &Self::Output { + &self.materializations_by_output[&output] } } diff --git a/data_plane/src/tests/test_utilities/engine_factories.rs b/data_plane/src/tests/test_utilities/engine_factories.rs index 3b0a540f..62a7621e 100644 --- a/data_plane/src/tests/test_utilities/engine_factories.rs +++ b/data_plane/src/tests/test_utilities/engine_factories.rs @@ -88,7 +88,7 @@ pub fn create_engine_single_pop_with_aggregated( .cloned() .collect(); - let mut materializations_by_policy_fingerprint = HashMap::new(); + let mut materializations_by_output = HashMap::new(); let agg_config = PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, @@ -118,13 +118,16 @@ pub fn create_engine_single_pop_with_aggregated( value_source_column: None, }; let agg_id = agg_config.policy_fp_u64(); - materializations_by_policy_fingerprint.insert(agg_id, agg_config); + materializations_by_output.insert(agg_id, agg_config); let installed_precompute_plan = Arc::new(InstalledPrecomputePlan { partitioning: Default::default(), raw_programs: Default::default(), precompute_plan: None, - materializations_by_policy_fingerprint, + materializations_by_output: materializations_by_output + .into_iter() + .map(|(id, cfg)| (asap_types::sds::StoredOutputId(id), cfg)) + .collect(), storage_backend: Default::default(), }); @@ -134,7 +137,7 @@ pub fn create_engine_single_pop_with_aggregated( // Insert data into SketchStore via the canonical helper (M2.3.6e). let agg_cfg = installed_precompute_plan - .get_aggregation_config(agg_id) + .get_aggregation_config(asap_types::sds::StoredOutputId(agg_id)) .cloned() .expect("agg config must be in installed_precompute_plan"); let timestamp = 1_000_000_u64; @@ -176,7 +179,7 @@ pub fn create_engine_dual_input( .cloned() .collect(); - let mut materializations_by_policy_fingerprint = HashMap::new(); + let mut materializations_by_output = HashMap::new(); // Value aggregation let value_agg_config = PrecomputeMaterialization { @@ -208,7 +211,7 @@ pub fn create_engine_dual_input( value_source_column: None, }; let value_id = value_agg_config.policy_fp_u64(); - materializations_by_policy_fingerprint.insert(value_id, value_agg_config); + materializations_by_output.insert(value_id, value_agg_config); // Keys aggregation let keys_agg_config = PrecomputeMaterialization { @@ -240,13 +243,16 @@ pub fn create_engine_dual_input( value_source_column: None, }; let keys_id = keys_agg_config.policy_fp_u64(); - materializations_by_policy_fingerprint.insert(keys_id, keys_agg_config); + materializations_by_output.insert(keys_id, keys_agg_config); let installed_precompute_plan = Arc::new(InstalledPrecomputePlan { partitioning: Default::default(), raw_programs: Default::default(), precompute_plan: None, - materializations_by_policy_fingerprint, + materializations_by_output: materializations_by_output + .into_iter() + .map(|(id, cfg)| (asap_types::sds::StoredOutputId(id), cfg)) + .collect(), storage_backend: Default::default(), }); @@ -255,11 +261,11 @@ pub fn create_engine_dual_input( let resolver = std::sync::Arc::new(SeriesIdResolver::new()); let agg_cfg_1 = installed_precompute_plan - .get_aggregation_config(value_id) + .get_aggregation_config(asap_types::sds::StoredOutputId(value_id)) .cloned() .expect("value agg config"); let agg_cfg_2 = installed_precompute_plan - .get_aggregation_config(keys_id) + .get_aggregation_config(asap_types::sds::StoredOutputId(keys_id)) .cloned() .expect("keys agg config"); let timestamp = 1_000_000_u64; @@ -308,7 +314,7 @@ pub fn create_engine_two_metrics( let labels_a: Vec = grouping_labels_a.iter().map(|s| s.to_string()).collect(); let labels_b: Vec = grouping_labels_b.iter().map(|s| s.to_string()).collect(); - let mut materializations_by_policy_fingerprint = HashMap::new(); + let mut materializations_by_output = HashMap::new(); let agg_config_a = PrecomputeMaterialization { stored_output_id: None, @@ -339,7 +345,7 @@ pub fn create_engine_two_metrics( value_source_column: None, }; let id_a = agg_config_a.policy_fp_u64(); - materializations_by_policy_fingerprint.insert(id_a, agg_config_a); + materializations_by_output.insert(id_a, agg_config_a); let agg_config_b = PrecomputeMaterialization { stored_output_id: None, @@ -370,13 +376,16 @@ pub fn create_engine_two_metrics( value_source_column: None, }; let id_b = agg_config_b.policy_fp_u64(); - materializations_by_policy_fingerprint.insert(id_b, agg_config_b); + materializations_by_output.insert(id_b, agg_config_b); let installed_precompute_plan = Arc::new(InstalledPrecomputePlan { partitioning: Default::default(), raw_programs: Default::default(), precompute_plan: None, - materializations_by_policy_fingerprint, + materializations_by_output: materializations_by_output + .into_iter() + .map(|(id, cfg)| (asap_types::sds::StoredOutputId(id), cfg)) + .collect(), storage_backend: Default::default(), }); @@ -384,11 +393,11 @@ pub fn create_engine_two_metrics( std::sync::Arc::new(crate::storage_engines::sketch_db::index::SketchStore::new()); let resolver = std::sync::Arc::new(SeriesIdResolver::new()); let agg_cfg_1 = installed_precompute_plan - .get_aggregation_config(id_a) + .get_aggregation_config(asap_types::sds::StoredOutputId(id_a)) .cloned() .expect("agg a"); let agg_cfg_2 = installed_precompute_plan - .get_aggregation_config(id_b) + .get_aggregation_config(asap_types::sds::StoredOutputId(id_b)) .cloned() .expect("agg b"); let timestamp = 1_000_000_u64; @@ -442,7 +451,7 @@ pub fn create_engine_three_metrics( let labels_b: Vec = grouping_labels_b.iter().map(|s| s.to_string()).collect(); let labels_c: Vec = grouping_labels_c.iter().map(|s| s.to_string()).collect(); - let mut materializations_by_policy_fingerprint = HashMap::new(); + let mut materializations_by_output = HashMap::new(); let mut ids: Vec = Vec::new(); for (agg_type, labels, metric) in [ @@ -480,14 +489,17 @@ pub fn create_engine_three_metrics( }; let id = cfg.policy_fp_u64(); ids.push(id); - materializations_by_policy_fingerprint.insert(id, cfg); + materializations_by_output.insert(id, cfg); } let installed_precompute_plan = Arc::new(InstalledPrecomputePlan { partitioning: Default::default(), raw_programs: Default::default(), precompute_plan: None, - materializations_by_policy_fingerprint, + materializations_by_output: materializations_by_output + .into_iter() + .map(|(id, cfg)| (asap_types::sds::StoredOutputId(id), cfg)) + .collect(), storage_backend: Default::default(), }); @@ -498,7 +510,7 @@ pub fn create_engine_three_metrics( .iter() .map(|id| { installed_precompute_plan - .get_aggregation_config(*id) + .get_aggregation_config(asap_types::sds::StoredOutputId(*id)) .cloned() .expect("agg present") }) @@ -535,7 +547,7 @@ pub fn create_engine_multi_timestamp( let grouping_label_strings: Vec = grouping_labels.iter().map(|s| s.to_string()).collect(); - let mut materializations_by_policy_fingerprint = HashMap::new(); + let mut materializations_by_output = HashMap::new(); let agg_config = PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, @@ -565,13 +577,16 @@ pub fn create_engine_multi_timestamp( value_source_column: None, }; let agg_id = agg_config.policy_fp_u64(); - materializations_by_policy_fingerprint.insert(agg_id, agg_config); + materializations_by_output.insert(agg_id, agg_config); let installed_precompute_plan = Arc::new(InstalledPrecomputePlan { partitioning: Default::default(), raw_programs: Default::default(), precompute_plan: None, - materializations_by_policy_fingerprint, + materializations_by_output: materializations_by_output + .into_iter() + .map(|(id, cfg)| (asap_types::sds::StoredOutputId(id), cfg)) + .collect(), storage_backend: Default::default(), }); @@ -579,7 +594,7 @@ pub fn create_engine_multi_timestamp( std::sync::Arc::new(crate::storage_engines::sketch_db::index::SketchStore::new()); let resolver = std::sync::Arc::new(SeriesIdResolver::new()); let agg_cfg = installed_precompute_plan - .get_aggregation_config(agg_id) + .get_aggregation_config(asap_types::sds::StoredOutputId(agg_id)) .cloned() .expect("agg"); for (timestamp, label_values_opt, acc) in data { @@ -613,7 +628,7 @@ pub fn create_engine_multi_timestamp_with_window( let grouping_label_strings: Vec = grouping_labels.iter().map(|s| s.to_string()).collect(); - let mut materializations_by_policy_fingerprint = HashMap::new(); + let mut materializations_by_output = HashMap::new(); let agg_config = PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, @@ -643,13 +658,16 @@ pub fn create_engine_multi_timestamp_with_window( value_source_column: None, }; let agg_id = agg_config.policy_fp_u64(); - materializations_by_policy_fingerprint.insert(agg_id, agg_config); + materializations_by_output.insert(agg_id, agg_config); let installed_precompute_plan = Arc::new(InstalledPrecomputePlan { partitioning: Default::default(), raw_programs: Default::default(), precompute_plan: None, - materializations_by_policy_fingerprint, + materializations_by_output: materializations_by_output + .into_iter() + .map(|(id, cfg)| (asap_types::sds::StoredOutputId(id), cfg)) + .collect(), storage_backend: Default::default(), }); @@ -657,7 +675,7 @@ pub fn create_engine_multi_timestamp_with_window( std::sync::Arc::new(crate::storage_engines::sketch_db::index::SketchStore::new()); let resolver = std::sync::Arc::new(SeriesIdResolver::new()); let agg_cfg = installed_precompute_plan - .get_aggregation_config(agg_id) + .get_aggregation_config(asap_types::sds::StoredOutputId(agg_id)) .cloned() .expect("agg"); for (timestamp, label_values_opt, acc) in data {