From 2c5e82eefcf2a67b0b8de8d3809f8d9b8d15a2ad Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 19:37:17 -0600 Subject: [PATCH 1/2] fix(storage): recover authoritative summary identity from disk --- .../storage_engines/sketch_db/index/mod.rs | 104 ++++++++++-- .../sketch_db/persistence/metadata.rs | 84 +++++++++- .../src/storage_engines/sketch_db/sds.rs | 31 +++- .../asapquery_compatibility_process_e2e.rs | 3 + .../tests/support/durable_summary_process.rs | 148 ++++++++++++++++++ .../continuous-summary-completeness.md | 2 +- 6 files changed, 352 insertions(+), 20 deletions(-) create mode 100644 data_plane/tests/support/durable_summary_process.rs 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 c90744b49..d0ceb0a7f 100644 --- a/data_plane/src/storage_engines/sketch_db/index/mod.rs +++ b/data_plane/src/storage_engines/sketch_db/index/mod.rs @@ -2744,6 +2744,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, @@ -2763,14 +2789,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 } @@ -2844,15 +2867,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( @@ -4289,6 +4324,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 = @@ -4307,6 +4385,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 { @@ -4329,9 +4410,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 0deee8a63..cffb47a5d 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..008e0c662 --- /dev/null +++ b/data_plane/tests/support/durable_summary_process.rs @@ -0,0 +1,148 @@ +//! 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..100 { + 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; + } + let metadata: Value = + serde_json::from_slice(&std::fs::read(&sidecar).expect("flushed metadata")).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 3f767f81b..5525ee462 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. From 247772fdf96a632f217cee80dffa3b1d154c73e2 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 19:45:34 -0600 Subject: [PATCH 2/2] test(storage): retain restart artifacts when durable flush times out --- data_plane/tests/support/durable_summary_process.rs | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/data_plane/tests/support/durable_summary_process.rs b/data_plane/tests/support/durable_summary_process.rs index 008e0c662..b37d39c80 100644 --- a/data_plane/tests/support/durable_summary_process.rs +++ b/data_plane/tests/support/durable_summary_process.rs @@ -106,7 +106,7 @@ async fn persisted_summary_restarts_without_live_reregistration() { .unwrap(); assert!(is_warm(&before), "{before}"); let sidecar = disk.join("sketch_index/sid_metadata.json"); - for _ in 0..100 { + for _ in 0..500 { let manifest_path = disk.join("sketch_index"); if sidecar.exists() && manifest_path.join("parts_manifest.log").exists() @@ -119,8 +119,14 @@ async fn persisted_summary_restarts_without_live_reregistration() { } tokio::time::sleep(Duration::from_millis(20)).await; } - let metadata: Value = - serde_json::from_slice(&std::fs::read(&sidecar).expect("flushed metadata")).unwrap(); + 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()