Skip to content
Open
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
2 changes: 1 addition & 1 deletion crates/asap_types/src/aggregation_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -321,7 +321,7 @@ impl PrecomputeMaterialization {

/// `PolicyFingerprint::as_u64()` — the u64-form handle used by the
/// policy-fingerprint-keyed call sites (e.g. `InstalledPrecomputePlan`'s
/// `HashMap<u64, PrecomputeMaterialization>` keys). **Always** equal to
/// `HashMap<StoredOutputId, PrecomputeMaterialization>` keys). **Always** equal to
/// `self.policy_fingerprint().as_u64()`. The value is content-
/// addressed identity, NOT a controller-allocated counter id.
pub fn policy_fp_u64(&self) -> u64 {
Expand Down
2 changes: 1 addition & 1 deletion crates/asap_types/src/policy_registry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
//!
//! Two `PrecomputeMaterialization`s that produce the same `PolicyFingerprint`
//! ARE the same policy. The registry treats this as a *deduplication*
//! invariant — if two distinct entries in the source `materializations_by_policy_fingerprint`
//! invariant — if two distinct entries in the source `materializations_by_output`
//! map produce the same fingerprint, the later one wins (last-write
//! semantics). In practice the source should never contain duplicates;
//! if it does, that's a control-plane bug worth surfacing in telemetry
Expand Down
12 changes: 10 additions & 2 deletions crates/asap_types/src/precompute_plan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -422,13 +422,21 @@ impl PrecomputePlan {

pub fn runtime_materializations(
&self,
) -> Result<HashMap<u64, crate::PrecomputeMaterialization>, PrecomputePlanError> {
) -> Result<
HashMap<crate::sds::StoredOutputId, crate::PrecomputeMaterialization>,
PrecomputePlanError,
> {
self.validate()?;
Ok(self
.materializations
.iter()
.cloned()
.map(|materialization| (materialization.policy_fp_u64(), materialization))
.map(|materialization| {
(
crate::sds::StoredOutputId(materialization.policy_fp_u64()),
materialization,
)
})
.collect())
}

Expand Down
5 changes: 5 additions & 0 deletions crates/asap_types/src/sds.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,11 @@ impl StoredOutputId {
self.0
}
}
impl From<u64> for StoredOutputId {
fn from(value: u64) -> Self {
Self(value)
}
}
impl From<crate::PolicyFingerprint> for StoredOutputId {
fn from(value: crate::PolicyFingerprint) -> Self {
Self(value.0)
Expand Down
4 changes: 2 additions & 2 deletions data_plane/src/drivers/ingest/otel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1405,7 +1405,7 @@ async fn route_modified_otlp_sketches_to_precompute(
// below so the query engine can answer per-item estimate(key)
// (the CMS/CountSketch FrequencyEstimate gate consults it).
let item_label_for_sid: Option<String> = {
snap.get_aggregation_config(policy_fp.as_u64())
snap.get_aggregation_config(policy_fp.into())
.or_else(|| {
snap.materializations()
.values()
Expand Down Expand Up @@ -4621,7 +4621,7 @@ mod sid_bucketing_tests {
let policy_fp = asap_types::PolicyFingerprint(cfg.policy_fp_u64());
let mut configs = HashMap::new();
configs.insert(cfg.policy_fp_u64(), cfg.clone());
let streaming = InstalledPrecomputePlan::new(configs);
let streaming = InstalledPrecomputePlan::from_raw_ids(configs);
let hot_reload = InstalledPrecomputePlanHandle::new(streaming);

let resolver = Arc::new(SeriesIdResolver::new());
Expand Down
19 changes: 10 additions & 9 deletions data_plane/src/drivers/ingest/prometheus_remote_write.rs
Original file line number Diff line number Diff line change
Expand Up @@ -435,7 +435,7 @@ impl PrometheusRemoteWriteReceiver {
continue;
};
let config = snapshot
.get_aggregation_config(policy_fp.as_u64())
.get_aggregation_config((*policy_fp).into())
.ok_or(RemoteWriteError::InactivePhysicalPlan)?;
let manager = crate::precompute_engine::window_manager::WindowManager::with_layout(
config.window_size,
Expand Down Expand Up @@ -987,7 +987,7 @@ mod tests {
capability_snapshot_id: "test".into(),
};
let configs = streaming
.materializations_by_policy_fingerprint
.materializations_by_output
.values()
.cloned()
.collect::<Vec<_>>();
Expand Down Expand Up @@ -1021,7 +1021,7 @@ mod tests {
producers: Vec::new(),
executable_dags: Default::default(),
materializations: streaming
.materializations_by_policy_fingerprint
.materializations_by_output
.values()
.cloned()
.collect(),
Expand Down Expand Up @@ -1094,7 +1094,8 @@ mod tests {
value_source_column: None,
};
let policy_fp = aggregation.policy_fp_u64();
let streaming = InstalledPrecomputePlan::new(HashMap::from([(policy_fp, aggregation)]));
let streaming =
InstalledPrecomputePlan::from_raw_ids(HashMap::from([(policy_fp, aggregation)]));
let (sender, receiver) = mpsc::channel(8);
let ingest = Arc::new(IngestState {
router: SeriesRouter::new(vec![sender]),
Expand Down Expand Up @@ -1132,7 +1133,7 @@ mod tests {
let mut config = snapshot.precompute_plan.materializations[0].clone();
config.population_key_encoding = asap_types::PopulationKeyEncoding::CanonicalLabelsV1;
config.partitioning = Some(asap_types::sds::PopulationPartitioning::Grouped);
let hot = physical_config(InstalledPrecomputePlan::new(HashMap::from([(
let hot = physical_config(InstalledPrecomputePlan::from_raw_ids(HashMap::from([(
config.policy_fp_u64(),
config.clone(),
)])));
Expand Down Expand Up @@ -1248,7 +1249,7 @@ mod tests {
assert_ne!(kll_fp, pooled_kll_fp);
let cms_fp = cms.policy_fingerprint();
let counter_fp = counter.policy_fingerprint();
let streaming = InstalledPrecomputePlan::new(HashMap::from([
let streaming = InstalledPrecomputePlan::from_raw_ids(HashMap::from([
(cms_fp.0, cms),
(counter_fp.0, counter),
(kll_fp.0, kll),
Expand Down Expand Up @@ -1652,18 +1653,18 @@ mod tests {
let policy = *ingest
.hot_reload_config
.snapshot()
.materializations_by_policy_fingerprint
.materializations_by_output
.keys()
.next()
.unwrap();
let binding = asap_types::query_plan::MaterializationBinding {
full_window_slide_ms: None,
materialization: asap_types::PolicyFingerprint(policy).into(),
materialization: policy,
stored_output_reference: ingest
.summary_store
.summary_catalog_snapshot()
.unwrap()
.output_reference(asap_types::PolicyFingerprint(policy).into())
.output_reference(policy)
.unwrap(),
output_grouping: asap_types::query_plan::PhysicalGrouping::Reduce(vec!["job".into()]),
item_labels: vec![],
Expand Down
9 changes: 5 additions & 4 deletions data_plane/src/drivers/query/servers/http.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3274,7 +3274,7 @@ mod tests {
marker_to_fp.insert(*marker, fp);
agg_map.insert(fp, cfg);
}
let installed_precompute_plan = Arc::new(InstalledPrecomputePlan::new(agg_map));
let installed_precompute_plan = Arc::new(InstalledPrecomputePlan::from_raw_ids(agg_map));
let hot_reload = InstalledPrecomputePlanHandle::from_arc(installed_precompute_plan.clone());
let query_engine = Arc::new(ASAPQueryEngine::new(15000));
let summary_store = Arc::new(crate::storage_engines::sketch_db::index::SketchStore::new());
Expand Down Expand Up @@ -5418,10 +5418,10 @@ pub fn validate_and_build_runtime_plan(
)
.map_err(|error| format!("DAG execution installation failed: {error}"))?;
let typed_fps: BTreeSet<_> = installed_precompute_plan
.materializations_by_policy_fingerprint
.materializations_by_output
.keys()
.copied()
.map(asap_types::PolicyFingerprint)
.map(asap_types::sds::StoredOutputId::fingerprint)
.collect();
request
.query_plan
Expand Down Expand Up @@ -6245,7 +6245,8 @@ async fn handle_post_backfill_job(
return (StatusCode::SERVICE_UNAVAILABLE, axum::Json(body)).into_response();
};
let snapshot = handle.snapshot();
let agg_cfg = match snapshot.get_aggregation_config(req.agg_id) {
let agg_cfg = match snapshot.get_aggregation_config(asap_types::sds::StoredOutputId(req.agg_id))
{
Some(c) => c.clone(),
None => {
let body = serde_json::json!({
Expand Down
2 changes: 1 addition & 1 deletion data_plane/src/precompute_engine/ingest_handler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -336,7 +336,7 @@ mod tests {

let mut map = std::collections::HashMap::new();
map.insert(agg_id, make_config(agg_id, metric));
let streaming = InstalledPrecomputePlan::new(map);
let streaming = InstalledPrecomputePlan::from_raw_ids(map);
let hot_reload =
crate::storage_engines::types::InstalledPrecomputePlanHandle::new(streaming.clone());

Expand Down
11 changes: 5 additions & 6 deletions data_plane/src/precompute_engine/output_sink.rs
Original file line number Diff line number Diff line change
Expand Up @@ -415,7 +415,7 @@ mod tests {
let agg_id = cfg.policy_fp_u64();
let mut configs = HashMap::new();
configs.insert(agg_id, cfg);
let streaming = InstalledPrecomputePlan::new(configs);
let streaming = InstalledPrecomputePlan::from_raw_ids(configs);
let hot_reload = InstalledPrecomputePlanHandle::new(streaming.clone());

let summary_store = Arc::new(SketchStore::new());
Expand Down Expand Up @@ -484,10 +484,9 @@ mod tests {
Arc::new(SeriesIdResolver::open(temporary.path().join("resolver.wal")).unwrap());
let sink = SketchStoreSink::new(
store.clone(),
InstalledPrecomputePlanHandle::new(InstalledPrecomputePlan::new(HashMap::from([(
fingerprint.0,
cfg,
)]))),
InstalledPrecomputePlanHandle::new(InstalledPrecomputePlan::from_raw_ids(
HashMap::from([(fingerprint.0, cfg)]),
)),
resolver,
);
let original_generation = Arc::new(catalog.reference().unwrap());
Expand Down Expand Up @@ -566,7 +565,7 @@ mod tests {
cfg.parameters
.insert("alpha".into(), serde_json::json!(0.01));
let policy_fp = cfg.policy_fp_u64();
let hot_reload = InstalledPrecomputePlanHandle::new(InstalledPrecomputePlan::new(
let hot_reload = InstalledPrecomputePlanHandle::new(InstalledPrecomputePlan::from_raw_ids(
HashMap::from([(policy_fp, cfg)]),
));
let summary_store = Arc::new(SketchStore::new());
Expand Down
2 changes: 1 addition & 1 deletion data_plane/src/precompute_engine/series_router.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ pub enum WorkerMessage {
sid: u64,
/// Source `PrecomputeMaterialization` fingerprint. Worker looks up its
/// `PrecomputeMaterialization` (window size, sketch kind/config, late
/// data policy, etc.) via `snap.get_aggregation_config(policy_fp.as_u64())`.
/// data policy, etc.) via `snap.get_aggregation_config(policy_fp.into())`.
policy_fp: PolicyFingerprint,
/// Grouping label values joined by semicolons (e.g. "constant").
/// Empty string if the aggregation has no grouping labels. Used
Expand Down
7 changes: 4 additions & 3 deletions data_plane/src/precompute_engine/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -513,7 +513,7 @@ impl Worker {
) -> Result<Option<&mut GroupState>, String> {
if !self.group_states.contains_key(&sid) {
let snap = self.hot_reload.snapshot();
let Some(cfg) = snap.get_aggregation_config(policy_fp.as_u64()) else {
let Some(cfg) = snap.get_aggregation_config(policy_fp.into()) else {
return Ok(None);
};
let program = snap.raw_programs.get(&policy_fp.as_u64()).cloned();
Expand Down Expand Up @@ -1081,7 +1081,7 @@ impl Worker {
let snap = self.hot_reload.snapshot();
let before = self.group_states.len();
self.group_states.retain(|&sid, gs| {
if snap.contains(gs.policy_fp.as_u64()) {
if snap.contains(gs.policy_fp.into()) {
return true; // policy still in config, keep
}
// Policy retired — keep only if there's residual data
Expand Down Expand Up @@ -1976,6 +1976,7 @@ mod tests {
use asap_physical_operators::summary_kernels::sum::SumAccumulator;
use asap_sketchlib::KllSketch;
use asap_types::enums::WindowKind;
use asap_types::sds::StoredOutputId;
use asap_types::AggregationType;

fn make_agg_config(
Expand Down Expand Up @@ -2086,7 +2087,7 @@ mod tests {
configs: HashMap<u64, PrecomputeMaterialization>,
) -> crate::storage_engines::types::InstalledPrecomputePlanHandle {
crate::storage_engines::types::InstalledPrecomputePlanHandle::new(
crate::storage_engines::types::InstalledPrecomputePlan::new(configs),
crate::storage_engines::types::InstalledPrecomputePlan::from_raw_ids(configs),
)
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1046,10 +1046,9 @@ mod tests {
name: "value".into(),
});
let hot = crate::storage_engines::types::InstalledPrecomputePlanHandle::from_arc(Arc::new(
crate::storage_engines::types::InstalledPrecomputePlan::new(HashMap::from([(
cfg.policy_fp_u64(),
cfg.clone(),
)])),
crate::storage_engines::types::InstalledPrecomputePlan::from_raw_ids(HashMap::from([
(cfg.policy_fp_u64(), cfg.clone()),
])),
));
let registry =
Arc::new(crate::storage_engines::sketch_db::backfill::BackfillRegistry::new());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -212,7 +212,7 @@ impl WindowProcessor for BackfillWindowProcessor {
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let snapshot = self.config.snapshot();
let config = snapshot
.get_aggregation_config(agg_id)
.get_aggregation_config(asap_types::sds::StoredOutputId(agg_id))
.cloned()
.ok_or_else(|| format!("agg_id {agg_id} not in current InstalledPrecomputePlan"))?;
let program = snapshot.raw_programs.get(&agg_id).cloned();
Expand Down Expand Up @@ -444,7 +444,7 @@ mod tests {
fn streaming_config_with(config: PrecomputeMaterialization) -> Arc<InstalledPrecomputePlan> {
let mut map = std::collections::HashMap::new();
map.insert(config.policy_fp_u64(), config);
Arc::new(InstalledPrecomputePlan::new(map))
Arc::new(InstalledPrecomputePlan::from_raw_ids(map))
}

#[tokio::test]
Expand Down
6 changes: 4 additions & 2 deletions data_plane/src/storage_engines/sketch_db/backfill/service.rs
Original file line number Diff line number Diff line change
Expand Up @@ -200,7 +200,9 @@ impl BackfillService {
);

let snapshot = self.config_source.snapshot();
let Some(materialization) = snapshot.get_aggregation_config(job.agg_id) else {
let Some(materialization) =
snapshot.get_aggregation_config(asap_types::sds::StoredOutputId(job.agg_id))
else {
self.registry
.mark_failed(job.job_id, "backfill materialization is not installed");
continue;
Expand Down Expand Up @@ -372,7 +374,7 @@ mod tests {
fn streaming_with(cfg: PrecomputeMaterialization) -> Arc<InstalledPrecomputePlan> {
let mut m = std::collections::HashMap::new();
m.insert(cfg.policy_fp_u64(), cfg);
Arc::new(InstalledPrecomputePlan::new(m))
Arc::new(InstalledPrecomputePlan::from_raw_ids(m))
}

async fn wait_for_status(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -244,7 +244,10 @@ mod tests {
id_to_fp.insert(id, fp);
map.insert(fp, cfg);
}
(Arc::new(InstalledPrecomputePlan::new(map)), id_to_fp)
(
Arc::new(InstalledPrecomputePlan::from_raw_ids(map)),
id_to_fp,
)
}

fn write_one(
Expand All @@ -263,7 +266,7 @@ mod tests {
asap_types::PolicyFingerprint(agg_id),
);
let agg_cfg = installed_precompute_plan
.get_aggregation_config(agg_id)
.get_aggregation_config(asap_types::sds::StoredOutputId(agg_id))
.unwrap();
// Test-scoped resolver — each call mints fresh. Production
// shares one resolver across all sinks; tests don't need that
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -326,7 +326,7 @@ mod tests {
for (i, c) in configs.into_iter().enumerate() {
map.insert(i as u64 + 1, c);
}
InstalledPrecomputePlan::new(map)
InstalledPrecomputePlan::from_raw_ids(map)
}

#[test]
Expand Down
Loading
Loading