From 30230ddfc4a1199519b1c5ef3d1206d3d80a10f9 Mon Sep 17 00:00:00 2001 From: zzylol Date: Tue, 29 Sep 2026 01:56:23 +0000 Subject: [PATCH] refactor: route precompute outputs on StoredOutputId, not a bare u64 The runtime's materialization index was keyed by a bare `u64` whose meaning lived only in the field name, `materializations_by_policy_fingerprint`. That name points at PolicyFingerprint, which its own module doc calls a "Legacy routing wrapper", superseded by an explicit `stored_output_id`. Every caller had to know which of the two identities that integer was. Key the index on StoredOutputId, the type whose doc says "Runtime routing uses this identity, never the semantic definition hash", and rename the field to `materializations_by_output`. `runtime_materializations`, the three accessors and the `Index` impl follow. No serialized bytes change. InstalledPrecomputePlan is `#[derive(Debug, Clone)]` with no Serialize at all, its own doc calling it an internal execution index derived from a validated plan. StoredOutputId and PolicyFingerprint are the same u64 behind two `#[serde(transparent)]` newtypes with conversions both ways, so even `state_schema_id`, which does reach a published document, formats the identical string. Fixtures keep holding raw ids. `from_raw_ids` is the one place that says so, rather than each test wrapping literals, and `From for StoredOutputId` lets a value already established as an output id convert without ceremony. This is step 1 of #785. It does not retire PolicyFingerprint: `from_config` still derives one from an explicit `stored_output_id`, so adoption changes how a fingerprint is computed rather than removing it. Retirement needs `stored_output_id` to become required, which is a separate change. Co-Authored-By: Claude Opus 5 (1M context) --- crates/asap_types/src/aggregation_config.rs | 2 +- crates/asap_types/src/policy_registry.rs | 2 +- crates/asap_types/src/precompute_plan.rs | 12 ++- crates/asap_types/src/sds.rs | 5 ++ data_plane/src/drivers/ingest/otel.rs | 4 +- .../drivers/ingest/prometheus_remote_write.rs | 19 ++--- data_plane/src/drivers/query/servers/http.rs | 9 ++- .../src/precompute_engine/ingest_handler.rs | 2 +- .../src/precompute_engine/output_sink.rs | 11 ++- .../src/precompute_engine/series_router.rs | 2 +- data_plane/src/precompute_engine/worker.rs | 7 +- .../accelerator.rs | 7 +- .../sketch_db/backfill/processor.rs | 4 +- .../sketch_db/backfill/service.rs | 6 +- .../sketch_db/lifecycle/eviction.rs | 7 +- .../sketch_db/lifecycle/reconcile.rs | 2 +- .../types/hot_reload_config.rs | 37 +++++----- .../types/installed_precompute_plan.rs | 55 ++++++++------ .../tests/test_utilities/engine_factories.rs | 74 ++++++++++++------- 19 files changed, 156 insertions(+), 111 deletions(-) diff --git a/crates/asap_types/src/aggregation_config.rs b/crates/asap_types/src/aggregation_config.rs index 24d804d2f..f256d74fd 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 16f007bb8..90138b139 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 f360ab4b9..05227b679 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 2e665f9c7..a6870ea20 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 7185b4792..81925b1a7 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 3816c0dbc..258c66cd6 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 1199e8867..5480a76f8 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 cf033c140..400797a7a 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 42f2562c1..eab97b320 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 6a29e4cf0..0334f5fb2 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 d957c0bf4..6e5230175 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 e50ae3f05..c8162ffb5 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 ceb6e9844..f72eaa481 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 35c9b1353..8e5639424 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 b3edf16f2..4788aaca0 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 c5bd2d504..d6310b46f 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 567852c80..0e37ac146 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 4ca677707..1ec253024 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 3b0a540fb..62a7621ec 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 {