Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
248 changes: 216 additions & 32 deletions data_plane/src/storage_engines/sketch_db/index/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -689,6 +689,8 @@ pub struct SketchStore {
/// has sealed epochs to persist) with retention-drop disabled (the
/// flush-then-evict loop is the memory bound).
persistence_read: RwLock<Option<Arc<PersistenceReadHandle>>>,
persistence_metadata: RwLock<Option<Arc<persistence::metadata::SidMetadataStore>>>,
removed_sids: RwLock<BTreeSet<u64>>,
/// 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`].
Expand Down Expand Up @@ -781,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);
Expand Down Expand Up @@ -2389,21 +2395,76 @@ impl SketchStore {
.collect()
}

fn metadata_record(&self, m: &SdsBinding) -> Option<persistence::metadata::SidMetaRecord> {
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<Arc<SketchInstanceMetadata>> {
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))
}

Expand All @@ -2416,10 +2477,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);
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))
}

Expand All @@ -2443,6 +2507,10 @@ 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()?;
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();
let removed = instances.remove(&sid);
Expand Down Expand Up @@ -2708,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<crate::storage_engines::sketch_db::index::persistence::Manifest>,
pub part_cache: crate::storage_engines::sketch_db::index::persistence::cache::PartCache,
Expand Down Expand Up @@ -2741,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,
Expand Down Expand Up @@ -2785,7 +2868,12 @@ impl SketchStore {
}),
);

let flusher = FlusherHandle::start(cfg, Arc::clone(&manifest), Arc::clone(self))?;
let flusher = FlusherHandle::start_with_metadata(
cfg,
Arc::clone(&manifest),
Arc::clone(self),
metadata_writer,
)?;

Ok(SketchIndexPersistence {
manifest,
Expand Down Expand Up @@ -2820,6 +2908,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;
Expand Down Expand Up @@ -2867,8 +2958,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() {
Expand Down Expand Up @@ -2947,27 +3038,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(
Expand Down Expand Up @@ -4479,6 +4550,119 @@ 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();
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());
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
.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"
);
recovered.register(meta_with_policy(803, fingerprint));
assert!(
recovered.instance(803).is_none(),
"removed SID reused after restart"
);
}

#[test]
fn observed_inventory_includes_durable_instances_after_restart() {
let snapshot: control_plane::physical::compiler::BackendLocalPlanningSnapshot =
Expand Down
21 changes: 18 additions & 3 deletions data_plane/src/storage_engines/sketch_db/persistence/flusher.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<super::metadata::SidMetadataStore>,
pub next_part_id: AtomicU64,
pub shutdown: AtomicBool,
/// Woken by the insert path when it hits `hard_cap_bytes` and by
Expand All @@ -66,6 +66,19 @@ impl FlusherHandle {
manifest: Arc<Manifest>,
source: Arc<S>,
) -> PersistResult<Self>
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<S>(
cfg: SketchStorePersistenceConfig,
manifest: Arc<Manifest>,
source: Arc<S>,
sid_metadata: Arc<super::metadata::SidMetadataStore>,
) -> PersistResult<Self>
where
S: EpochSource + 'static,
{
Expand All @@ -79,8 +92,6 @@ impl FlusherHandle {
.map(|m| m + 1)
.unwrap_or(1);

let sid_metadata = super::metadata::SidMetadataStore::new(&cfg.disk_path);

let shared = Arc::new(FlusherShared {
cfg: cfg.clone(),
manifest: manifest.clone(),
Expand All @@ -104,6 +115,10 @@ impl FlusherHandle {
})
}

pub(crate) fn metadata_store(&self) -> Arc<super::metadata::SidMetadataStore> {
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) {
Expand Down
Loading
Loading