From d708e568429cb9e6e0d07d390ec28ec125e2e293 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 20:02:54 -0600 Subject: [PATCH 1/2] fix(storage): persist summary lifecycle tombstones --- .../storage_engines/sketch_db/index/mod.rs | 189 +++++++++++++++--- .../sketch_db/persistence/flusher.rs | 8 +- .../sketch_db/persistence/metadata.rs | 46 ++++- .../continuous-summary-completeness.md | 2 + 4 files changed, 211 insertions(+), 34 deletions(-) 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 f480e00ab..8e41c7f4c 100644 --- a/data_plane/src/storage_engines/sketch_db/index/mod.rs +++ b/data_plane/src/storage_engines/sketch_db/index/mod.rs @@ -689,6 +689,7 @@ pub struct SketchStore { /// has sealed epochs to persist) with retention-drop disabled (the /// flush-then-evict loop is the memory bound). persistence_read: RwLock>>, + persistence_metadata: RwLock>>, /// Seal cadence in distinct windows, applied to every per-sid /// `SidStoreData` once persistence is enabled. `0` (the default) /// disables cadence sealing. Set by [`Self::enable_persistence_mode`]. @@ -2389,21 +2390,76 @@ impl SketchStore { .collect() } + fn metadata_record(&self, m: &SdsBinding) -> Option { + 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); + } + record.retired_at_ms = m.retired_at_ms; + record.expires_at_ms = m.expires_at_ms; + Some(record) + } + + fn persist_lifecycle(&self, instance: &SdsBinding, removed: bool) -> Result<(), String> { + let Some(writer) = self.persistence_metadata.read().unwrap().clone() else { + return Ok(()); + }; + // A retired definition can disappear from the desired catalog before + // its stored instances are collected. Preserve its persisted provenance + // rather than rebinding it to the new catalog generation. + let mut record = match self.metadata_record(instance) { + Some(record) => record, + None => writer + .load() + .map_err(|error| error.to_string())? + .into_iter() + .find(|record| record.sid == instance.sid) + .ok_or("cannot persist lifecycle without catalog identity")?, + }; + record.retired_at_ms = instance.retired_at_ms; + record.expires_at_ms = instance.expires_at_ms; + record.removed = removed; + writer.upsert_all(&[record]).map_err(|error| { + tracing::error!(sid = instance.sid, %error, "durable summary lifecycle publication failed"); + error.to_string() + }) + } + /// Force `sid` into `Retired` status, scheduling expiry /// `retention` from now. Idempotent — re-retiring a Retired or /// Expired sid is a no-op and returns the unchanged metadata. - /// Returns `None` if the sid is unknown. + /// Returns `None` if the sid is unknown or durable lifecycle publication fails. pub fn force_retire( &self, sid: u64, retention: Duration, ) -> Option> { + let _mutation = self.begin_state_mutation(); let mut map = self.instances.write().ok()?; let instance = map.get_mut(&sid)?; - let meta = Arc::make_mut(&mut instance.metadata); + let mut next = instance.clone(); + let meta = Arc::make_mut(&mut next.metadata); if matches!(meta.status(), AggStatus::Active) { meta.retire(retention); } + self.persist_lifecycle(&next, false).ok()?; + *instance = next; Some(Arc::clone(&instance.metadata)) } @@ -2416,10 +2472,13 @@ impl SketchStore { let _mutation = self.begin_state_mutation(); let mut map = self.instances.write().ok()?; let instance = map.get_mut(&sid)?; - let meta = Arc::make_mut(&mut instance.metadata); + let mut next = instance.clone(); + let meta = Arc::make_mut(&mut next.metadata); let now = now_ms(); meta.retired_at_ms = Some(now); meta.expires_at_ms = Some(now); + self.persist_lifecycle(&next, false).ok()?; + *instance = next; Some(Arc::clone(&instance.metadata)) } @@ -2443,6 +2502,9 @@ impl SketchStore { let removed = { // Fixed lock order: instances → policy_to_series_ids → metric_to_series_ids. let mut instances = self.instances.write().ok()?; + if let Some(instance) = instances.get(&sid) { + self.persist_lifecycle(instance, true).ok()?; + } let mut policy_idx = self.policy_to_series_ids.write().unwrap(); let mut metric_idx = self.metric_to_series_ids.write().unwrap(); let removed = instances.remove(&sid); @@ -2786,6 +2848,7 @@ impl SketchStore { ); let flusher = FlusherHandle::start(cfg, Arc::clone(&manifest), Arc::clone(self))?; + *self.persistence_metadata.write().unwrap() = Some(flusher.metadata_store()); Ok(SketchIndexPersistence { manifest, @@ -2820,6 +2883,9 @@ impl SketchStore { let mut registered = 0usize; for rec in records { + if rec.removed || rec.expires_at_ms.is_some_and(|expiry| expiry <= now_ms()) { + continue; + } // Don't clobber a live-registered instance. if self.instance(rec.sid).is_some() { continue; @@ -2867,8 +2933,8 @@ impl SketchStore { agg_kind, accuracy, first_seen_unix_ms: rec.first_seen_unix_ms, - retired_at_ms: None, - expires_at_ms: None, + retired_at_ms: rec.retired_at_ms, + expires_at_ms: rec.expires_at_ms, policy_fp, }); if self.instance(rec.sid).is_some() { @@ -2947,27 +3013,7 @@ impl crate::storage_engines::sketch_db::index::persistence::EpochSource for Sket { let g = self.instances.read().ok()?; let m = g.get(&sid)?; - 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) + self.metadata_record(m) } fn snapshot_sealed_epoch( @@ -4479,6 +4525,97 @@ mod tests { assert!(store.series_ids_for_policy(fingerprint).is_empty()); } + #[test] + fn failed_durable_lifecycle_write_preserves_live_instance() { + let store = SketchStore::new(); + store.register(meta(800)); + let directory = tempfile::tempdir().unwrap(); + let writer = Arc::new(persistence::metadata::SidMetadataStore::new( + directory.path(), + )); + // A directory in place of the sidecar causes the real writer to fail. + std::fs::create_dir(writer.path()).unwrap(); + *store.persistence_metadata.write().unwrap() = Some(writer); + assert!(store.force_retire(800, Duration::from_secs(60)).is_none()); + assert!(store.force_expire(800).is_none()); + assert!(store.remove_instance(800).is_none()); + let instance = store.instance(800).unwrap(); + assert!(instance.retired_at_ms.is_none()); + assert!(instance.expires_at_ms.is_none()); + } + + #[test] + fn durable_lifecycle_is_not_resurrected_by_restart_or_a_stale_flush() { + 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 directory = tempfile::tempdir().unwrap(); + let disk = directory.path().to_path_buf(); + let expected_retirement; + { + let store = Arc::new(SketchStore::new()); + store + .install_summary_catalog(Arc::new(plan.summary_catalog.clone())) + .unwrap(); + for sid in [801, 802, 803] { + store.register(meta_with_policy(sid, fingerprint)); + } + let mut persistence = store.start_persistence(durable_cfg(disk.clone())).unwrap(); + for sid in [801, 802, 803] { + for pane in 0..4 { + store.append_sample( + sid, + BTreeMap::new(), + (pane * 30_000, (pane + 1) * 30_000), + sample(1), + ); + } + } + assert!(wait_until( + || !persistence.manifest.live_parts().is_empty(), + Duration::from_secs(5) + )); + let stale: Vec<_> = { + let instances = store.instances.read().unwrap(); + [801, 802, 803] + .iter() + .map(|sid| store.metadata_record(&instances[sid]).unwrap()) + .collect() + }; + expected_retirement = store.force_retire(801, Duration::from_secs(3600)).unwrap(); + assert!(store.force_expire(802).is_some()); + assert!(store.remove_instance(803).is_some()); + // This models a flush that captured metadata before the lifecycle + // operation and reaches the shared writer afterward. + persistence + .flusher + .metadata_store() + .upsert_all(&stale) + .unwrap(); + persistence.shutdown(); + } + let recovered = Arc::new(SketchStore::new()); + recovered + .install_summary_catalog(Arc::new(plan.summary_catalog)) + .unwrap(); + let _persistence = recovered.start_persistence(durable_cfg(disk)).unwrap(); + let retired = recovered.instance(801).unwrap(); + assert_eq!(retired.retired_at_ms, expected_retirement.retired_at_ms); + assert_eq!(retired.expires_at_ms, expected_retirement.expires_at_ms); + assert!( + recovered.instance(802).is_none(), + "expired state resurrected" + ); + assert!( + recovered.instance(803).is_none(), + "removed state resurrected" + ); + } + #[test] fn observed_inventory_includes_durable_instances_after_restart() { let snapshot: control_plane::physical::compiler::BackendLocalPlanningSnapshot = diff --git a/data_plane/src/storage_engines/sketch_db/persistence/flusher.rs b/data_plane/src/storage_engines/sketch_db/persistence/flusher.rs index 1fc87b64f..0ff66a16a 100644 --- a/data_plane/src/storage_engines/sketch_db/persistence/flusher.rs +++ b/data_plane/src/storage_engines/sketch_db/persistence/flusher.rs @@ -42,7 +42,7 @@ pub(crate) struct FlusherShared { /// Per-sid metadata sidecar — upserted on every flush so recovery can /// re-register disk-resident sids as queryable instances. See /// [`super::metadata`]. - pub sid_metadata: super::metadata::SidMetadataStore, + pub sid_metadata: Arc, pub next_part_id: AtomicU64, pub shutdown: AtomicBool, /// Woken by the insert path when it hits `hard_cap_bytes` and by @@ -79,7 +79,7 @@ impl FlusherHandle { .map(|m| m + 1) .unwrap_or(1); - let sid_metadata = super::metadata::SidMetadataStore::new(&cfg.disk_path); + let sid_metadata = Arc::new(super::metadata::SidMetadataStore::new(&cfg.disk_path)); let shared = Arc::new(FlusherShared { cfg: cfg.clone(), @@ -104,6 +104,10 @@ impl FlusherHandle { }) } + pub(crate) fn metadata_store(&self) -> Arc { + Arc::clone(&self.inner.sid_metadata) + } + /// Signal shutdown and wait for the thread to finish its current /// tick. Safe to call multiple times; subsequent calls are no-ops. pub fn shutdown(&mut self) { 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 a3b7b2eaa..51e6e683d 100644 --- a/data_plane/src/storage_engines/sketch_db/persistence/metadata.rs +++ b/data_plane/src/storage_engines/sketch_db/persistence/metadata.rs @@ -271,6 +271,12 @@ pub struct SidMetaRecord { pub group_by_keys: Vec, agg_kind: AggKindRec, pub first_seen_unix_ms: i64, + #[serde(default)] + pub retired_at_ms: Option, + #[serde(default)] + pub expires_at_ms: Option, + #[serde(default)] + pub removed: bool, } impl SidMetaRecord { @@ -292,6 +298,9 @@ impl SidMetaRecord { group_by_keys, agg_kind: agg_kind.into(), first_seen_unix_ms, + retired_at_ms: None, + expires_at_ms: None, + removed: false, } } @@ -353,6 +362,12 @@ struct SidBindingRec { summary_descriptor_id: String, data_descriptor_id: String, first_seen_unix_ms: i64, + #[serde(default)] + retired_at_ms: Option, + #[serde(default)] + expires_at_ms: Option, + #[serde(default)] + removed: bool, } /// Version-3 normalized sidecar with authoritative catalog provenance. Descriptors appear once and SeriesId bindings hold @@ -420,6 +435,9 @@ impl SdsSidecar { summary_descriptor_id: summary_id, data_descriptor_id: data_id, first_seen_unix_ms: record.first_seen_unix_ms, + retired_at_ms: record.retired_at_ms, + expires_at_ms: record.expires_at_ms, + removed: record.removed, }, ); } @@ -470,6 +488,9 @@ impl SdsSidecar { group_by_keys: data.group_by_keys.clone(), agg_kind: operator.with_population_filter(&data.population_filter_canonical), first_seen_unix_ms: binding.first_seen_unix_ms, + retired_at_ms: binding.retired_at_ms, + expires_at_ms: binding.expires_at_ms, + removed: binding.removed, }) }) .collect() @@ -483,6 +504,7 @@ impl SdsSidecar { #[derive(Debug)] pub struct SidMetadataStore { path: PathBuf, + writer: std::sync::Mutex<()>, } impl SidMetadataStore { @@ -491,6 +513,7 @@ impl SidMetadataStore { pub fn new(disk_path: &Path) -> Self { Self { path: disk_path.join(SERIES_ID_METADATA_FILE), + writer: std::sync::Mutex::new(()), } } @@ -566,6 +589,9 @@ impl SidMetadataStore { /// disk (last write wins per sid). Atomic via tmp + rename + dir /// fsync, matching the manifest's durability discipline. pub fn upsert_all(&self, records: &[SidMetaRecord]) -> PersistResult<()> { + let _writer = self.writer.lock().map_err(|_| { + PersistError::Io(std::io::Error::other("summary metadata writer poisoned")) + })?; if records.is_empty() { return Ok(()); } @@ -577,12 +603,20 @@ impl SidMetadataStore { let mut changed = false; for r in records { let key = r.sid.to_string(); - match map.get(&key) { - Some(existing) if existing == r => {} - _ => { - map.insert(key, r.clone()); - changed = true; - } + let mut next = r.clone(); + if let Some(existing) = map.get(&key) { + // Lifecycle is monotone for a SeriesId. An older flush snapshot + // must not resurrect a retired or removed persisted instance. + next.removed |= existing.removed; + next.retired_at_ms = existing.retired_at_ms.or(next.retired_at_ms); + next.expires_at_ms = match (existing.expires_at_ms, next.expires_at_ms) { + (Some(a), Some(b)) => Some(a.min(b)), + (a, b) => a.or(b), + }; + } + if map.get(&key) != Some(&next) { + map.insert(key, next); + changed = true; } } if !changed { diff --git a/docs/design_docs/continuous-summary-completeness.md b/docs/design_docs/continuous-summary-completeness.md index 30ab3dd13..34549b8df 100644 --- a/docs/design_docs/continuous-summary-completeness.md +++ b/docs/design_docs/continuous-summary-completeness.md @@ -11,3 +11,5 @@ 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, 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. + +Durable SID bindings also preserve retirement and expiry timestamps and a removal tombstone. Lifecycle changes publish through the same serialized metadata writer as the flusher before changing in-memory visibility. A stale flush snapshot cannot clear those fields. Recovery leaves removed and expired instances unregistered and preserves a still-retired instance's expiry deadline. Tombstones remain until durable state is explicitly reclaimed; this change does not claim automatic tombstone garbage collection or cross-generation reactivation. From 523441f9f0386b31691ab5d3585fa39d0ba03714 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 20:19:05 -0600 Subject: [PATCH 2/2] fix(storage): keep lifecycle identity monotone during recovery --- .../storage_engines/sketch_db/index/mod.rs | 61 ++++++++++++++++--- .../sketch_db/persistence/flusher.rs | 15 ++++- 2 files changed, 67 insertions(+), 9 deletions(-) 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 8e41c7f4c..6ac76aef4 100644 --- a/data_plane/src/storage_engines/sketch_db/index/mod.rs +++ b/data_plane/src/storage_engines/sketch_db/index/mod.rs @@ -690,6 +690,7 @@ pub struct SketchStore { /// flush-then-evict loop is the memory bound). persistence_read: RwLock>>, persistence_metadata: RwLock>>, + removed_sids: RwLock>, /// Seal cadence in distinct windows, applied to every per-sid /// `SidStoreData` once persistence is enabled. `0` (the default) /// disables cadence sealing. Set by [`Self::enable_persistence_mode`]. @@ -782,6 +783,10 @@ impl SketchStore { .map(str::to_owned); // Fixed lock order: instances → policy_to_series_ids → metric_to_series_ids. let mut instances = self.instances.write().unwrap(); + if self.removed_sids.read().unwrap().contains(&sid) { + tracing::warn!(sid, "rejecting reuse of a removed summary instance ID"); + return; + } let mut policy_idx = self.policy_to_series_ids.write().unwrap(); let mut metric_idx = self.metric_to_series_ids.write().unwrap(); instances.insert(sid, instance); @@ -2475,8 +2480,8 @@ impl SketchStore { let mut next = instance.clone(); let meta = Arc::make_mut(&mut next.metadata); let now = now_ms(); - meta.retired_at_ms = Some(now); - meta.expires_at_ms = Some(now); + meta.retired_at_ms = Some(meta.retired_at_ms.unwrap_or(now).min(now)); + meta.expires_at_ms = Some(meta.expires_at_ms.unwrap_or(now).min(now)); self.persist_lifecycle(&next, false).ok()?; *instance = next; Some(Arc::clone(&instance.metadata)) @@ -2504,6 +2509,7 @@ impl SketchStore { let mut instances = self.instances.write().ok()?; if let Some(instance) = instances.get(&sid) { self.persist_lifecycle(instance, true).ok()?; + self.removed_sids.write().ok()?.insert(sid); } let mut policy_idx = self.policy_to_series_ids.write().unwrap(); let mut metric_idx = self.metric_to_series_ids.write().unwrap(); @@ -2770,9 +2776,8 @@ impl SketchStore { /// impl on `SketchStore` and writes parts under `disk_path/parts/`. /// /// Drop or call [`Self::shutdown`] to stop the flusher cleanly. The -/// `part_cache` field is exposed so the query path can be wired up to -/// read-back from disk in a subsequent sub-PR; today it sits idle -/// because the in-memory `query_range` doesn't yet consult it. +/// `part_cache` backs disk reads; `query_range` combines durable parts with +/// live in-memory state through the installed persistence read handle. pub struct SketchIndexPersistence { pub manifest: Arc, pub part_cache: crate::storage_engines::sketch_db::index::persistence::cache::PartCache, @@ -2803,6 +2808,22 @@ impl SketchStore { cache::PartCache, flusher::FlusherHandle, recovery, Manifest, }; + // Publish the writer before recovery or any background work so a + // concurrent lifecycle operation cannot succeed without persistence. + let metadata_writer = + Arc::new(persistence::metadata::SidMetadataStore::new(&cfg.disk_path)); + { + // Registration and lifecycle changes take this lock first too. + let _instances = self.instances.write().unwrap(); + *self.persistence_metadata.write().unwrap() = Some(Arc::clone(&metadata_writer)); + self.removed_sids.write().unwrap().extend( + metadata_writer + .load()? + .into_iter() + .filter(|record| record.removed) + .map(|record| record.sid), + ); + } let (_loaded_manifest, report) = recovery::recover(&cfg.disk_path)?; tracing::info!( live = report.live_parts, @@ -2847,8 +2868,12 @@ impl SketchStore { }), ); - let flusher = FlusherHandle::start(cfg, Arc::clone(&manifest), Arc::clone(self))?; - *self.persistence_metadata.write().unwrap() = Some(flusher.metadata_store()); + let flusher = FlusherHandle::start_with_metadata( + cfg, + Arc::clone(&manifest), + Arc::clone(self), + metadata_writer, + )?; Ok(SketchIndexPersistence { manifest, @@ -4525,6 +4550,18 @@ mod tests { assert!(store.series_ids_for_policy(fingerprint).is_empty()); } + #[test] + fn force_expire_never_extends_existing_lifecycle_deadlines() { + let store = SketchStore::new(); + let mut metadata = meta(799); + metadata.retired_at_ms = Some(1); + metadata.expires_at_ms = Some(2); + store.register(metadata); + let expired = store.force_expire(799).unwrap(); + assert_eq!(expired.retired_at_ms, Some(1)); + assert_eq!(expired.expires_at_ms, Some(2)); + } + #[test] fn failed_durable_lifecycle_write_preserves_live_instance() { let store = SketchStore::new(); @@ -4589,6 +4626,11 @@ mod tests { expected_retirement = store.force_retire(801, Duration::from_secs(3600)).unwrap(); assert!(store.force_expire(802).is_some()); assert!(store.remove_instance(803).is_some()); + store.register(meta_with_policy(803, fingerprint)); + assert!( + store.instance(803).is_none(), + "removed SID reused before restart" + ); // This models a flush that captured metadata before the lifecycle // operation and reaches the shared writer afterward. persistence @@ -4614,6 +4656,11 @@ mod tests { recovered.instance(803).is_none(), "removed state resurrected" ); + recovered.register(meta_with_policy(803, fingerprint)); + assert!( + recovered.instance(803).is_none(), + "removed SID reused after restart" + ); } #[test] diff --git a/data_plane/src/storage_engines/sketch_db/persistence/flusher.rs b/data_plane/src/storage_engines/sketch_db/persistence/flusher.rs index 0ff66a16a..c2aad589b 100644 --- a/data_plane/src/storage_engines/sketch_db/persistence/flusher.rs +++ b/data_plane/src/storage_engines/sketch_db/persistence/flusher.rs @@ -66,6 +66,19 @@ impl FlusherHandle { manifest: Arc, source: Arc, ) -> PersistResult + where + S: EpochSource + 'static, + { + let metadata = Arc::new(super::metadata::SidMetadataStore::new(&cfg.disk_path)); + Self::start_with_metadata(cfg, manifest, source, metadata) + } + + pub(crate) fn start_with_metadata( + cfg: SketchStorePersistenceConfig, + manifest: Arc, + source: Arc, + sid_metadata: Arc, + ) -> PersistResult where S: EpochSource + 'static, { @@ -79,8 +92,6 @@ impl FlusherHandle { .map(|m| m + 1) .unwrap_or(1); - let sid_metadata = Arc::new(super::metadata::SidMetadataStore::new(&cfg.disk_path)); - let shared = Arc::new(FlusherShared { cfg: cfg.clone(), manifest: manifest.clone(),