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
104 changes: 93 additions & 11 deletions data_plane/src/storage_engines/sketch_db/index/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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
}
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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 =
Expand All @@ -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 {
Expand All @@ -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();
Expand Down
84 changes: 79 additions & 5 deletions data_plane/src/storage_engines/sketch_db/persistence/metadata.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<asap_types::sds::SummaryDefinitionId>,
#[serde(default)]
pub catalog_generation: Option<std::sync::Arc<asap_types::sds::CatalogGeneration>>,
pub metric_name: String,
/// Label KEY set, sorted (a `Vec` so the JSON stays compact; the
/// store side rebuilds the `BTreeSet`).
Expand All @@ -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(),
Expand Down Expand Up @@ -339,16 +346,22 @@ struct DataDescriptorRec {
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
struct SidBindingRec {
sid: u64,
#[serde(default)]
summary_definition_id: Option<asap_types::sds::SummaryDefinitionId>,
#[serde(default)]
catalog_generation_sha256: Option<String>,
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<String, std::sync::Arc<asap_types::sds::CatalogGeneration>>,
summary_descriptors: HashMap<String, AggKindRec>,
data_descriptors: HashMap<String, DataDescriptorRec>,
bindings: HashMap<String, SidBindingRec>,
Expand All @@ -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(),
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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) => {
Expand Down Expand Up @@ -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();
Expand All @@ -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
Expand Down Expand Up @@ -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);
}

Expand Down
31 changes: 28 additions & 3 deletions data_plane/src/storage_engines/sketch_db/sds.rs
Original file line number Diff line number Diff line change
Expand Up @@ -101,7 +101,12 @@ pub struct SummaryDescriptorRegistry {
// registry must not turn retired materializations into a permanent leak.
summaries: RwLock<HashMap<SummaryDescriptorId, Weak<SummaryDescriptor>>>,
data: RwLock<HashMap<DataDescriptorId, Weak<DataDescriptor>>>,
authoritative_catalog: RwLock<Option<Arc<asap_types::summary_catalog::SummaryCatalog>>>,
authoritative_catalog: RwLock<
Option<(
Arc<asap_types::summary_catalog::SummaryCatalog>,
Arc<asap_types::sds::CatalogGeneration>,
)>,
>,
}

impl SummaryDescriptorRegistry {
Expand All @@ -110,18 +115,38 @@ impl SummaryDescriptorRegistry {
catalog: Arc<asap_types::summary_catalog::SummaryCatalog>,
) -> 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<Arc<asap_types::summary_catalog::SummaryCatalog>> {
self.authoritative_catalog
.read()
.unwrap()
.as_ref()
.map(|(catalog, _)| Arc::clone(catalog))
}

pub fn authoritative_snapshot(
&self,
) -> Option<(
Arc<asap_types::summary_catalog::SummaryCatalog>,
Arc<asap_types::sds::CatalogGeneration>,
)> {
self.authoritative_catalog.read().unwrap().clone()
}

pub fn bind(&self, metadata: SketchInstanceMetadata) -> Result<SdsBinding, String> {
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(
Expand Down
3 changes: 3 additions & 0 deletions data_plane/tests/asapquery_compatibility_process_e2e.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Loading
Loading