diff --git a/data_plane/src/storage_engines/sketch_db/index/mod.rs b/data_plane/src/storage_engines/sketch_db/index/mod.rs index 8a6074f5c..f480e00ab 100644 --- a/data_plane/src/storage_engines/sketch_db/index/mod.rs +++ b/data_plane/src/storage_engines/sketch_db/index/mod.rs @@ -2824,6 +2824,32 @@ impl SketchStore { if self.instance(rec.sid).is_some() { continue; } + let catalog = self.descriptors.authoritative_snapshot(); + let policy_fp = match ( + &rec.summary_definition_id, + &rec.catalog_generation, + &catalog, + ) { + (Some(definition), Some(generation), Some((catalog, installed_generation))) => { + if generation != installed_generation + || !catalog.materializations.contains_key(definition) + { + tracing::warn!( + sid = rec.sid, + "persisted summary catalog provenance differs; leaving state unbound" + ); + continue; + } + definition.fingerprint() + } + // Legacy deployments without an authoritative plan retain their + // legacy path. Never promote such state into a catalog binding. + (None, None, None) => PolicyFingerprint::UNSET, + _ => { + tracing::warn!(sid = rec.sid, "persisted summary has no matching authoritative identity; leaving state unbound"); + continue; + } + }; let Some(agg_kind) = rec.agg_kind() else { tracing::warn!( sid = rec.sid, @@ -2843,14 +2869,11 @@ impl SketchStore { first_seen_unix_ms: rec.first_seen_unix_ms, retired_at_ms: None, expires_at_ms: None, - // The sidecar doesn't carry the policy fingerprint; the - // recovered sid is reachable through the - // `instances_matching(metric, gbk)` walk regardless (the - // policy_fp reverse index is an optimization, not a - // correctness requirement for the query path). - policy_fp: PolicyFingerprint::UNSET, + policy_fp, }); - registered += 1; + if self.instance(rec.sid).is_some() { + registered += 1; + } } registered } @@ -2924,15 +2947,27 @@ impl crate::storage_engines::sketch_db::index::persistence::EpochSource for Sket { let g = self.instances.read().ok()?; let m = g.get(&sid)?; - Some( + let mut record = crate::storage_engines::sketch_db::index::persistence::metadata::SidMetaRecord::new( m.sid, m.metric_name.clone(), m.group_by_keys.iter().cloned().collect(), &m.agg_kind, m.first_seen_unix_ms, - ), - ) + ); + if !m.policy_fp.is_unset() { + let (catalog, generation) = self.descriptors.authoritative_snapshot()?; + let definition = SummaryDefinitionId::from(m.policy_fp); + let identity = catalog.materializations.get(&definition)?; + if identity.summary_descriptor_id != *m.summary_descriptor.id() + || identity.data_descriptor_id != *m.data_descriptor.id() + { + return None; + } + record.summary_definition_id = Some(definition); + record.catalog_generation = Some(generation); + } + Some(record) } fn snapshot_sealed_epoch( @@ -4401,6 +4436,49 @@ mod tests { drop(p2); } + #[test] + fn catalog_recovery_keeps_legacy_and_foreign_generation_state_unbound() { + use crate::storage_engines::sketch_db::index::persistence::metadata::{ + SidMetaRecord, SidMetadataStore, + }; + let snapshot: control_plane::physical::compiler::BackendLocalPlanningSnapshot = + serde_json::from_str(include_str!( + "../../../../../docs/examples/asapquery-compatibility-demo-snapshot.json" + )) + .unwrap(); + let plan = snapshot.compile().unwrap(); + let fingerprint = plan.precompute_plan.materializations[0].policy_fingerprint(); + let metadata = meta_with_policy(507, fingerprint); + let record = SidMetaRecord::new( + metadata.sid, + metadata.metric_name.clone(), + metadata.group_by_keys.iter().cloned().collect(), + &metadata.agg_kind, + 0, + ); + let tmp = tempfile::tempdir().unwrap(); + let sidecar = SidMetadataStore::new(tmp.path()); + sidecar.upsert_all(&[record.clone()]).unwrap(); + let store = SketchStore::new(); + store + .install_summary_catalog(Arc::new(plan.summary_catalog.clone())) + .unwrap(); + assert_eq!(store.register_recovered_disk_series(tmp.path()), 0); + assert!(store.instance(507).is_none()); + let mut foreign = record; + foreign.summary_definition_id = Some(fingerprint.into()); + let reference = plan.summary_catalog.reference().unwrap(); + foreign.catalog_generation = Some(Arc::new(CatalogGeneration { + schema_version: reference.schema_version, + plan_id: reference.plan_id, + plan_version: reference.plan_version + 1, + snapshot_sha256: reference.snapshot_sha256, + })); + sidecar.upsert_all(&[foreign]).unwrap(); + assert_eq!(store.register_recovered_disk_series(tmp.path()), 0); + assert!(store.series_ids_for_policy(fingerprint).is_empty()); + } + #[test] fn observed_inventory_includes_durable_instances_after_restart() { let snapshot: control_plane::physical::compiler::BackendLocalPlanningSnapshot = @@ -4419,6 +4497,9 @@ mod tests { { let store = Arc::new(SketchStore::new()); + store + .install_summary_catalog(Arc::new(plan.summary_catalog.clone())) + .unwrap(); store.register(metadata.clone()); let mut persistence = store.start_persistence(durable_cfg(disk.clone())).unwrap(); for index in 0..4u64 { @@ -4441,9 +4522,10 @@ mod tests { recovered .install_summary_catalog(Arc::new(plan.summary_catalog)) .unwrap(); - recovered.register(metadata); let persistence = recovered.start_persistence(durable_cfg(disk)).unwrap(); assert!(!persistence.manifest.live_parts().is_empty()); + assert_eq!(recovered.series_ids_for_policy(fingerprint), vec![506]); + assert!(!recovered.query_range(506, 0, 120_000).is_empty()); let inventory = recovered .observed_summary_inventory("backend-a", "store-a", &producers, 1, 100) .unwrap(); diff --git a/data_plane/src/storage_engines/sketch_db/persistence/metadata.rs b/data_plane/src/storage_engines/sketch_db/persistence/metadata.rs index c489e16c4..a3b7b2eaa 100644 --- a/data_plane/src/storage_engines/sketch_db/persistence/metadata.rs +++ b/data_plane/src/storage_engines/sketch_db/persistence/metadata.rs @@ -260,6 +260,11 @@ impl AggKindRec { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub struct SidMetaRecord { pub sid: u64, + /// Authoritative identity and provenance, absent on legacy sidecars. + #[serde(default)] + pub summary_definition_id: Option, + #[serde(default)] + pub catalog_generation: Option>, pub metric_name: String, /// Label KEY set, sorted (a `Vec` so the JSON stays compact; the /// store side rebuilds the `BTreeSet`). @@ -281,6 +286,8 @@ impl SidMetaRecord { ) -> Self { Self { sid, + summary_definition_id: None, + catalog_generation: None, metric_name, group_by_keys, agg_kind: agg_kind.into(), @@ -339,16 +346,22 @@ struct DataDescriptorRec { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] struct SidBindingRec { sid: u64, + #[serde(default)] + summary_definition_id: Option, + #[serde(default)] + catalog_generation_sha256: Option, summary_descriptor_id: String, data_descriptor_id: String, first_seen_unix_ms: i64, } -/// Version-2 normalized sidecar. Descriptors appear once and SeriesId bindings hold +/// Version-3 normalized sidecar with authoritative catalog provenance. Descriptors appear once and SeriesId bindings hold /// foreign keys, mirroring the in-memory SDS registry. #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] struct SdsSidecar { schema_version: u32, + #[serde(default)] + catalog_generations: HashMap>, summary_descriptors: HashMap, data_descriptors: HashMap, bindings: HashMap, @@ -359,7 +372,8 @@ impl SdsSidecar { use crate::storage_engines::sketch_db::sds::{data_descriptor_id, summary_descriptor_id}; let mut sidecar = Self { - schema_version: 2, + schema_version: 3, + catalog_generations: HashMap::new(), summary_descriptors: HashMap::new(), data_descriptors: HashMap::new(), bindings: HashMap::new(), @@ -389,10 +403,20 @@ impl SdsSidecar { population_filter_canonical: filter, group_by_keys: record.group_by_keys, }); + let generation_sha256 = record.catalog_generation.map(|generation| { + let digest = generation.snapshot_sha256.clone(); + sidecar + .catalog_generations + .entry(digest.clone()) + .or_insert(generation); + digest + }); sidecar.bindings.insert( record.sid.to_string(), SidBindingRec { sid: record.sid, + summary_definition_id: record.summary_definition_id, + catalog_generation_sha256: generation_sha256, summary_descriptor_id: summary_id, data_descriptor_id: data_id, first_seen_unix_ms: record.first_seen_unix_ms, @@ -426,6 +450,22 @@ impl SdsSidecar { })?; Ok(SidMetaRecord { sid: binding.sid, + summary_definition_id: binding.summary_definition_id, + catalog_generation: binding + .catalog_generation_sha256 + .as_ref() + .map(|digest| { + self.catalog_generations + .get(digest) + .cloned() + .ok_or_else(|| { + PersistError::Format(format!( + "SeriesId {} references missing catalog generation", + binding.sid + )) + }) + }) + .transpose()?, metric_name: data.metric_name.clone(), group_by_keys: data.group_by_keys.clone(), agg_kind: operator.with_population_filter(&data.population_filter_canonical), @@ -484,7 +524,10 @@ impl SidMetadataStore { return Ok(Vec::new()); } }; - if value.get("schema_version").and_then(|v| v.as_u64()) == Some(2) { + if matches!( + value.get("schema_version").and_then(|v| v.as_u64()), + Some(2 | 3) + ) { let sidecar: SdsSidecar = match serde_json::from_value(value) { Ok(sidecar) => sidecar, Err(error) => { @@ -615,6 +658,37 @@ mod tests { assert!(s.load().unwrap().is_empty()); } + #[test] + fn authoritative_bindings_share_one_persisted_catalog_generation() { + let directory = tempfile::tempdir().unwrap(); + let store = SidMetadataStore::new(directory.path()); + let generation = std::sync::Arc::new(asap_types::sds::CatalogGeneration { + schema_version: 1, + plan_id: 7, + plan_version: 2, + snapshot_sha256: "catalog".into(), + }); + let mut first = sketch_meta(1); + first.summary_definition_id = Some(asap_types::PolicyFingerprint(7).into()); + first.catalog_generation = Some(std::sync::Arc::clone(&generation)); + let mut second = first.clone(); + second.sid = 2; + store.upsert_all(&[first, second]).unwrap(); + let json: serde_json::Value = + serde_json::from_slice(&std::fs::read(store.path()).unwrap()).unwrap(); + assert_eq!(json["catalog_generations"].as_object().unwrap().len(), 1); + let records = store.load().unwrap(); + assert_eq!(records.len(), 2); + assert!(std::sync::Arc::ptr_eq( + records[0].catalog_generation.as_ref().unwrap(), + records[1].catalog_generation.as_ref().unwrap() + )); + assert_eq!( + records[0].summary_definition_id, + Some(asap_types::PolicyFingerprint(7).into()) + ); + } + #[test] fn upsert_then_load_round_trips() { let tmp = TempDir::new().unwrap(); @@ -629,7 +703,7 @@ mod tests { let persisted: serde_json::Value = serde_json::from_slice(&std::fs::read(s.path()).unwrap()).unwrap(); - assert_eq!(persisted["schema_version"], 2); + assert_eq!(persisted["schema_version"], 3); assert_eq!( persisted["summary_descriptors"].as_object().unwrap().len(), 2 @@ -689,7 +763,7 @@ mod tests { store.upsert_all(&[exact_meta(2)]).unwrap(); let persisted: serde_json::Value = serde_json::from_slice(&std::fs::read(store.path()).unwrap()).unwrap(); - assert_eq!(persisted["schema_version"], 2); + assert_eq!(persisted["schema_version"], 3); assert_eq!(store.load().unwrap().len(), 2); } diff --git a/data_plane/src/storage_engines/sketch_db/sds.rs b/data_plane/src/storage_engines/sketch_db/sds.rs index e54cac850..beecf8f14 100644 --- a/data_plane/src/storage_engines/sketch_db/sds.rs +++ b/data_plane/src/storage_engines/sketch_db/sds.rs @@ -101,7 +101,12 @@ pub struct SummaryDescriptorRegistry { // registry must not turn retired materializations into a permanent leak. summaries: RwLock>>, data: RwLock>>, - authoritative_catalog: RwLock>>, + authoritative_catalog: RwLock< + Option<( + Arc, + Arc, + )>, + >, } impl SummaryDescriptorRegistry { @@ -110,18 +115,38 @@ impl SummaryDescriptorRegistry { catalog: Arc, ) -> Result<(), asap_types::summary_catalog::SummaryCatalogError> { catalog.validate()?; - *self.authoritative_catalog.write().unwrap() = Some(catalog); + let reference = catalog.reference()?; + let generation = Arc::new(asap_types::sds::CatalogGeneration { + schema_version: reference.schema_version, + plan_id: reference.plan_id, + plan_version: reference.plan_version, + snapshot_sha256: reference.snapshot_sha256, + }); + *self.authoritative_catalog.write().unwrap() = Some((catalog, generation)); Ok(()) } pub fn authoritative_catalog( &self, ) -> Option> { + self.authoritative_catalog + .read() + .unwrap() + .as_ref() + .map(|(catalog, _)| Arc::clone(catalog)) + } + + pub fn authoritative_snapshot( + &self, + ) -> Option<( + Arc, + Arc, + )> { self.authoritative_catalog.read().unwrap().clone() } pub fn bind(&self, metadata: SketchInstanceMetadata) -> Result { - let authoritative = self.authoritative_catalog.read().unwrap().clone(); + let authoritative = self.authoritative_catalog(); let configured = if let Some(catalog) = authoritative.as_ref() { if metadata.policy_fp.is_unset() { return Err( diff --git a/data_plane/tests/asapquery_compatibility_process_e2e.rs b/data_plane/tests/asapquery_compatibility_process_e2e.rs index 24482f17c..d5f536111 100644 --- a/data_plane/tests/asapquery_compatibility_process_e2e.rs +++ b/data_plane/tests/asapquery_compatibility_process_e2e.rs @@ -21,6 +21,9 @@ use tokio::sync::Mutex; #[path = "support/erp_planning_process.rs"] mod erp_planning_process; +#[path = "support/durable_summary_process.rs"] +mod durable_summary_process; + struct ChildGuard(Child); impl Drop for ChildGuard { diff --git a/data_plane/tests/support/durable_summary_process.rs b/data_plane/tests/support/durable_summary_process.rs new file mode 100644 index 000000000..b37d39c80 --- /dev/null +++ b/data_plane/tests/support/durable_summary_process.rs @@ -0,0 +1,154 @@ +//! Authoritative catalog identity survives a real production process restart. +use super::*; + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn persisted_summary_restarts_without_live_reregistration() { + let mut fixture: Value = serde_json::from_str(include_str!( + "../../../docs/examples/asapquery-compatibility-demo-snapshot.json" + )) + .unwrap(); + let query = fixture["query_workload"]["repeating_queries"][2]["query"] + .as_str() + .unwrap() + .to_owned(); + fixture["query_workload"]["repeating_queries"] = + serde_json::json!([fixture["query_workload"]["repeating_queries"][2].clone()]); + let snapshot: control_plane::physical::compiler::BackendLocalPlanningSnapshot = + serde_json::from_value(fixture).unwrap(); + let plan = snapshot.compile().unwrap(); + let install = data_plane::drivers::query::servers::http::PhysicalPlanInstallRequest { + summary_catalog: plan.summary_catalog, + collector_plans: plan.collector_plans, + precompute_plan: plan.precompute_plan, + transmission_plan: plan.transmission_plan, + query_plan: plan.query_plan, + storage_routing: None, + adaptation_evidence: vec![], + }; + let directory = tempfile::tempdir().unwrap(); + let artifact = directory.path().join("plan.json"); + let bootstrap = directory.path().join("bootstrap.json"); + let disk = directory.path().join("disk"); + std::fs::create_dir_all(&disk).unwrap(); + std::fs::write(&artifact, serde_json::to_vec(&install).unwrap()).unwrap(); + std::fs::write(&bootstrap, b"{\"aggregations\":[]}").unwrap(); + let spawn = |port: u16| { + ChildGuard( + Command::new(env!("CARGO_BIN_EXE_data_plane")) + .arg("--physical-plan") + .arg(&artifact) + .arg("--streaming-config") + .arg(&bootstrap) + .arg("--http-port") + .arg(port.to_string()) + .arg("--output-dir") + .arg(directory.path()) + .arg("--enable-remote-write") + .arg("--persistence-enabled") + .arg("--persistence-dir") + .arg(&disk) + .args([ + "--persistence-memory-limit-mb", + "1", + "--persistence-hot-window-secs", + "1", + "--persistence-delete-older-than-secs", + "0", + "--persistence-seal-window-count", + "1", + "--persistence-flush-interval-ms", + "10", + "--precompute-allowed-lateness-ms", + "0", + "--precompute-flush-interval-ms", + "25", + ]) + .stdout(Stdio::null()) + .stderr(Stdio::inherit()) + .spawn() + .unwrap(), + ) + }; + let client = reqwest::Client::new(); + let port = unused_port(); + let base = format!("http://127.0.0.1:{port}"); + let mut first = spawn(port); + wait_until_ready(&client, &format!("{base}/api/v1/health"), &mut first.0).await; + assert_eq!( + remote_write( + &client, + &base, + &WriteRequest { + timeseries: vec![series( + "asap_demo_gauge", + &[ + (1000, 1.0), + (2000, 2.0), + (3000, 3.0), + (4000, 4.0), + (5000, 5.0) + ] + )], + } + ) + .await, + 204 + ); + drain_precompute(&client, &base).await; + let before: Value = client + .get(format!("{base}/api/v1/query")) + .query(&[("query", query.as_str()), ("time", "5")]) + .send() + .await + .unwrap() + .json() + .await + .unwrap(); + assert!(is_warm(&before), "{before}"); + let sidecar = disk.join("sketch_index/sid_metadata.json"); + for _ in 0..500 { + let manifest_path = disk.join("sketch_index"); + if sidecar.exists() + && manifest_path.join("parts_manifest.log").exists() + && data_plane::storage_engines::sketch_db::index::persistence::Manifest::open_or_init( + &manifest_path, + ) + .is_ok_and(|manifest| manifest.live_parts().iter().any(|part| part.max_ts >= 5000)) + { + break; + } + tokio::time::sleep(Duration::from_millis(20)).await; + } + if !sidecar.exists() { + let retained = directory.keep(); + panic!( + "missing flushed metadata; retained process artifacts at {}", + retained.display() + ); + } + let metadata: Value = serde_json::from_slice(&std::fs::read(&sidecar).unwrap()).unwrap(); + assert!(metadata["bindings"] + .as_object() + .unwrap() + .values() + .all(|binding| !binding["summary_definition_id"].is_null() + && !binding["catalog_generation_sha256"].is_null())); + drop(first); + let port = unused_port(); + let base = format!("http://127.0.0.1:{port}"); + let mut second = spawn(port); + wait_until_ready(&client, &format!("{base}/api/v1/health"), &mut second.0).await; + // No Remote Write, register call, or installation endpoint after restart. + let after: Value = client + .get(format!("{base}/api/v1/query")) + .query(&[("query", query.as_str()), ("time", "5")]) + .send() + .await + .unwrap() + .json() + .await + .unwrap(); + assert!(is_warm(&after), "{after}"); + assert_eq!(after["data"]["result"], before["data"]["result"]); + assert_eq!(after["data"]["result"][0]["value"][1], "15"); +} diff --git a/docs/design_docs/continuous-summary-completeness.md b/docs/design_docs/continuous-summary-completeness.md index 4be73e780..30ab3dd13 100644 --- a/docs/design_docs/continuous-summary-completeness.md +++ b/docs/design_docs/continuous-summary-completeness.md @@ -10,4 +10,4 @@ Finite drain closes input, waits for all workers and certifies that every accept Admission metadata has bounded coordinate, pending-revision and byte budgets. Completed receipts expire with configured materialization retention; pending work and the published prefixes of pending admissions remain protected. Already admitted slow-worker outputs can finish behind another worker's maintenance replay frontier. Unsolicited expired input and expired untagged replay remain rejected. -This inventory is not durable and does not establish exactly-once execution across crashes. Startup installs the authoritative catalog before persistence recovery, but persisted series metadata still needs a separate migration to preserve summary definition identity and catalog provenance. General multi-input maintenance transforms and durable producer watermarks remain separate work. +This inventory is not durable and does not establish exactly-once execution across crashes. Startup installs the authoritative catalog before persistence recovery, and version-3 persisted series metadata preserves summary definition identity and catalog provenance. Recovery requires the same catalog generation and leaves incompatible or legacy records unbound under an authoritative catalog. Reuse across changed catalog generations requires a separate explicit compatibility decision. General multi-input maintenance transforms and durable producer watermarks remain separate work.