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
4 changes: 2 additions & 2 deletions control_plane/src/opamp/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1018,7 +1018,7 @@ mod tests {
},
materializations: vec![crate::physical::compiler::CollectorMaterialization {
query_id: "q".into(),
materialization: asap_types::PolicyFingerprint(1),
materialization: asap_types::PolicyFingerprint(1).into(),
metric: "requests".into(),
algorithm: "hll".into(),
parameters: serde_json::json!({"precision": 14}),
Expand All @@ -1039,7 +1039,7 @@ mod tests {
},
}],
transmission_rules: vec![crate::physical::compiler::TransmissionRule {
materialization: asap_types::PolicyFingerprint(1),
materialization: asap_types::PolicyFingerprint(1).into(),
producer_id: collector_id.into(),
schema_id: "asap-query-backend.v1:summary-state:v1:1".into(),
mode: crate::physical::compiler::TransmissionMode::Full,
Expand Down
56 changes: 28 additions & 28 deletions control_plane/src/physical/compiler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -263,7 +263,7 @@ pub struct PlanEnvelope {
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct CollectorMaterialization {
pub query_id: String,
pub materialization: asap_types::PolicyFingerprint,
pub materialization: asap_types::sds::SummaryDefinitionId,
pub metric: String,
pub algorithm: String,
pub parameters: Value,
Expand Down Expand Up @@ -412,7 +412,7 @@ impl TryFrom<&SummaryFamilyType> for StateFamilyContract {
pub struct StateSchemaContract {
pub schema_id: String,
pub schema_version: u32,
pub materialization: asap_types::PolicyFingerprint,
pub materialization: asap_types::sds::SummaryDefinitionId,
pub family: StateFamilyContract,
pub source: Source,
pub value_column: planner_types::pre_asap::ColumnRef,
Expand Down Expand Up @@ -442,7 +442,7 @@ pub struct StateWindowContract {
pub struct ProducerContract {
pub producer_id: String,
pub collector_id: String,
pub materialization: asap_types::PolicyFingerprint,
pub materialization: asap_types::sds::SummaryDefinitionId,
pub schema_id: String,
}

Expand Down Expand Up @@ -501,7 +501,7 @@ impl PrecomputePlan {
Ok(StateSchemaContract {
schema_id: state_schema_id(fingerprint),
schema_version: 1,
materialization: fingerprint,
materialization: fingerprint.into(),
family,
source,
value_column,
Expand Down Expand Up @@ -615,7 +615,7 @@ impl PrecomputePlan {
materialization.policy_fp_u64(),
));
}
if !materializations.insert(materialization.policy_fingerprint()) {
if !materializations.insert(materialization.policy_fingerprint().into()) {
return Err(PrecomputePlanError::DuplicateMaterialization(
materialization.policy_fp_u64(),
));
Expand Down Expand Up @@ -650,13 +650,16 @@ impl PrecomputePlan {
let materialization = self
.materializations
.iter()
.find(|candidate| candidate.policy_fingerprint() == schema.materialization)
.find(|candidate| {
candidate.policy_fingerprint() == schema.materialization.fingerprint()
})
.ok_or(PrecomputePlanError::SchemaSetMismatch)?;
let accumulator = materialization
.accumulator_spec()
.map_err(|_| PrecomputePlanError::UnsupportedFamily(schema.materialization.0))?;
let family = StateFamilyContract::try_from(&accumulator.family)
.map_err(|_| PrecomputePlanError::UnsupportedFamily(schema.materialization.0))?;
let accumulator = materialization.accumulator_spec().map_err(|_| {
PrecomputePlanError::UnsupportedFamily(schema.materialization.as_u64())
})?;
let family = StateFamilyContract::try_from(&accumulator.family).map_err(|_| {
PrecomputePlanError::UnsupportedFamily(schema.materialization.as_u64())
})?;
let source = materialization.table_name.as_ref().map_or_else(
|| Source::TimeSeries {
metric: materialization.metric.clone(),
Expand All @@ -670,7 +673,7 @@ impl PrecomputePlan {
.clone()
.map(planner_types::pre_asap::ColumnRef::Named)
.unwrap_or(planner_types::pre_asap::ColumnRef::SampleValue);
if schema.schema_id != state_schema_id(schema.materialization)
if schema.schema_id != state_schema_id(schema.materialization.fingerprint())
|| schema.family != family
|| schema.source != source
|| schema.value_column != value_column
Expand Down Expand Up @@ -720,7 +723,7 @@ impl PrecomputePlan {
}
if self.ingest.require_registered_producer {
if let Some(missing) = materializations.difference(&produced).next() {
return Err(PrecomputePlanError::MissingProducer(missing.0));
return Err(PrecomputePlanError::MissingProducer(missing.as_u64()));
}
}
Ok(())
Expand All @@ -744,7 +747,7 @@ pub struct PhysicalPlan {

#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct MaterializationLifecycleEstimate {
pub materialization: asap_types::PolicyFingerprint,
pub materialization: asap_types::sds::SummaryDefinitionId,
pub consumer_query_ids: Vec<String>,
pub window_implementation_id: String,
pub horizon_seconds: f64,
Expand Down Expand Up @@ -778,7 +781,7 @@ pub struct FrameIdentityContract {
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct TransmissionRule {
pub materialization: asap_types::PolicyFingerprint,
pub materialization: asap_types::sds::SummaryDefinitionId,
pub producer_id: String,
pub schema_id: String,
pub mode: TransmissionMode,
Expand Down Expand Up @@ -981,7 +984,7 @@ pub struct RuntimeRulePolicy {
pub struct RuntimeAdaptationEvidence {
pub plan_id: u64,
pub plan_version: u64,
pub materialization: asap_types::PolicyFingerprint,
pub materialization: asap_types::sds::SummaryDefinitionId,
pub producer_id: String,
pub schema_id: String,
pub producer_version: String,
Expand Down Expand Up @@ -1017,7 +1020,7 @@ pub struct SummaryFrameIdentity {
pub plan_id: u64,
pub plan_version: u64,
pub backend_compat: String,
pub materialization: asap_types::PolicyFingerprint,
pub materialization: asap_types::sds::SummaryDefinitionId,
/// Canonical producer-side identity for one concrete retained-label group.
pub series_identity: String,
pub schema_id: String,
Expand Down Expand Up @@ -1062,7 +1065,7 @@ pub enum TransmissionPlanError {
fn validate_catalog_projection(
reference: Option<&super::summary_catalog::SummaryCatalogReference>,
envelope: &PlanEnvelope,
materializations: impl IntoIterator<Item = asap_types::PolicyFingerprint>,
materializations: impl IntoIterator<Item = asap_types::sds::SummaryDefinitionId>,
catalog: &super::summary_catalog::SummaryCatalog,
) -> Result<(), TransmissionPlanError> {
let expected = catalog
Expand All @@ -1077,13 +1080,10 @@ fn validate_catalog_projection(
));
}
for id in materializations {
if !catalog
.materializations
.contains_key(&asap_types::sds::SummaryDefinitionId::from(id))
{
if !catalog.materializations.contains_key(&id) {
return Err(TransmissionPlanError::Catalog(format!(
"unknown materialization {}",
id.0
id.as_u64()
)));
}
}
Expand Down Expand Up @@ -1143,10 +1143,10 @@ impl TransmissionPlan {
let materialization = precompute
.materializations
.iter()
.find(|m| m.policy_fingerprint() == producer.materialization)
.find(|m| m.policy_fingerprint() == producer.materialization.fingerprint())
.expect("validated PrecomputePlan materialization binding");
let runtime_policy = runtime_policies
.get(&producer.materialization)
.get(&producer.materialization.fingerprint())
.cloned()
.unwrap_or_default();
let mode = if runtime_policy.delta.is_some() {
Expand Down Expand Up @@ -2371,7 +2371,7 @@ impl PhysicalCompiler {
lifecycle_estimates
.entry(materialization)
.or_insert_with(|| MaterializationLifecycleEstimate {
materialization,
materialization: materialization.into(),
consumer_query_ids,
window_implementation_id: window_implementation.implementation_id.clone(),
horizon_seconds: query.lifecycle.horizon_seconds,
Expand Down Expand Up @@ -2414,7 +2414,7 @@ impl PhysicalCompiler {
compiled_materializations.push(runtime_materialization.clone());
let collector_materialization = CollectorMaterialization {
query_id: format!("state-{}", materialization.0),
materialization,
materialization: materialization.into(),
metric: metric.clone(),
algorithm: physical_algorithm,
parameters: selected.parameters,
Expand Down Expand Up @@ -5399,7 +5399,7 @@ mod tests {
assert!(validate_catalog_projection(
Some(&catalog.reference().unwrap()),
&bundle.envelope,
[asap_types::PolicyFingerprint(u64::MAX)],
[asap_types::PolicyFingerprint(u64::MAX).into()],
catalog
)
.is_err());
Expand Down
2 changes: 1 addition & 1 deletion control_plane/src/physical/precompute_contract.rs
Original file line number Diff line number Diff line change
Expand Up @@ -83,7 +83,7 @@ impl PrecomputePlan {
let schema = self
.schemas
.iter()
.find(|s| s.materialization == id.fingerprint())
.find(|s| s.materialization == id)
.ok_or(PrecomputePlanError::SchemaSetMismatch)?;
let family = config
.accumulator_spec()
Expand Down
2 changes: 1 addition & 1 deletion control_plane/src/physical/workload_cost.rs
Original file line number Diff line number Diff line change
Expand Up @@ -171,7 +171,7 @@ pub fn manifest(
.precompute_plan
.materializations
.iter()
.find(|m| m.policy_fingerprint() == schema.materialization)
.find(|m| m.policy_fingerprint() == schema.materialization.fingerprint())
.ok_or_else(|| invalid("state has no physical implementation"))?;
let identity = json!({"schema": schema, "location": location, "physical": physical,
"window_implementation": plan.lifecycle_estimates.iter().find(|e| e.materialization == schema.materialization).map(|e| &e.window_implementation_id)});
Expand Down
2 changes: 1 addition & 1 deletion crates/asap_types/src/sds.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ macro_rules! descriptor_id {
};
}
/// Semantic materialization reference. Wire-compatible with PolicyFingerprint,
/// but distinct from descriptor IDs and runtime instance/SID identity.
/// but distinct from descriptor IDs and concrete [`SummaryInstanceId`] identity.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(transparent)]
pub struct SummaryDefinitionId(pub crate::PolicyFingerprint);
Expand Down
17 changes: 10 additions & 7 deletions data_plane/src/drivers/ingest/otel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1265,7 +1265,7 @@ async fn route_modified_otlp_sketches_to_precompute(
&dp.attrs.keys().cloned().collect(),
)
});
if observed_policy != frame.materialization {
if observed_policy != frame.materialization.fingerprint() {
unreachable!(
"materialization changed after successful request preflight"
);
Expand Down Expand Up @@ -1637,7 +1637,7 @@ async fn route_modified_otlp_sketches_to_precompute(
debug!(
plan_id = frame.plan_id,
plan_version = frame.plan_version,
materialization = frame.materialization.0,
materialization = frame.materialization.as_u64(),
producer = %frame.producer_id,
producer_epoch = %frame.producer_epoch,
sequence = frame.sequence,
Expand All @@ -1652,7 +1652,7 @@ async fn route_modified_otlp_sketches_to_precompute(
warn!(
plan_id = frame.plan_id,
plan_version = frame.plan_version,
materialization = frame.materialization.0,
materialization = frame.materialization.as_u64(),
producer = %frame.producer_id,
producer_epoch = %frame.producer_epoch,
sequence = frame.sequence,
Expand Down Expand Up @@ -2237,10 +2237,10 @@ fn preflight_summary_frames(
&dp.attrs.keys().cloned().collect(),
)
};
if observed != frame.materialization {
if observed != frame.materialization.fingerprint() {
return Err(format!(
"summary frame for {metric_name} declares materialization {} but active schema resolves {}",
frame.materialization.0, observed.0
frame.materialization.as_u64(), observed.0
));
}
Ok(frame)
Expand Down Expand Up @@ -2392,7 +2392,7 @@ fn take_summary_frame_identity(
plan_id,
plan_version,
backend_compat,
materialization,
materialization: materialization.into(),
series_identity,
schema_id,
producer_id,
Expand Down Expand Up @@ -4695,7 +4695,10 @@ mod sid_bucketing_tests {
let frame = take_summary_frame_identity(&mut attrs, 100, 200).expect("valid identity");
assert_eq!(frame.plan_id, 42);
assert_eq!(frame.plan_version, 3);
assert_eq!(frame.materialization, asap_types::PolicyFingerprint(99));
assert_eq!(
frame.materialization,
asap_types::PolicyFingerprint(99).into()
);
assert_eq!(attrs, HashMap::from([("service".into(), "api".into())]));
}
}
4 changes: 2 additions & 2 deletions data_plane/src/precompute_engine/frame_lineage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ pub enum FrameLineageError {
struct FrameLineageKey {
plan_id: u64,
plan_version: u64,
materialization: asap_types::PolicyFingerprint,
materialization: asap_types::sds::SummaryDefinitionId,
series_identity: String,
producer_id: String,
producer_epoch: String,
Expand Down Expand Up @@ -199,7 +199,7 @@ mod tests {
plan_id: 7,
plan_version: 3,
backend_compat: "asap-query-backend.v1".into(),
materialization: asap_types::PolicyFingerprint(41),
materialization: asap_types::PolicyFingerprint(41).into(),
series_identity: "service=checkout,zone=a".into(),
schema_id: "schema-41".into(),
producer_id: "edge-a".into(),
Expand Down
2 changes: 1 addition & 1 deletion data_plane/src/storage_engines/sketch_db/index/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3130,7 +3130,7 @@ mod tests {
plan_id: 7,
plan_version: 2,
backend_compat: "asap-query-backend.v1".into(),
materialization: PolicyFingerprint(41),
materialization: PolicyFingerprint(41).into(),
series_identity: "service=checkout,zone=a".into(),
schema_id: "schema-41".into(),
producer_id: "edge-a".into(),
Expand Down
Loading